跳到主要内容

Kafka 生产者

插件用途

将转发组产生的变量、设备、报警和插件事件发布到 Apache Kafka Topic。

转发组范围、触发、分批、缓存和启停见数据转发。本页只说明 Kafka 连接、安全认证、Topic 模板、上传模板和专用调试。

功能入口

进入“开发配置 → 数据转发”,按以下顺序配置:

  1. 新增或编辑转发组,先保存变量范围、触发模式、定时间隔、在线过滤和批处理策略。
  2. 新增目标,选择“Kafka 生产者”,填写目标名称、启用状态、日志级别和启动超时。
  3. 打开“目标属性”,依次填写 Kafka 地址、安全认证、消息配置、数据与脚本以及缓存容量。
  4. 保存并启用转发组和目标,再从“目标调试”验证发布结果。

目标基本信息

参数默认值如何配置
所属转发组-必须选择已保存的转发组;目标只接收该组范围内的变量。
目标名称-必填,同一转发组内不能重复。建议包含环境和 Kafka 集群名称。
启用开启关闭时 Kafka 生产者不会启动。
日志级别Info排查连接或发布问题时临时使用 Debug,完成后恢复。
启动超时60页面可填 13600 秒。

目标属性

连接配置

参数默认值如何配置
服务地址127.0.0.1:9092填 Kafka Bootstrap Server。多个节点用逗号分隔,例如 kafka-a:9092,kafka-b:9092。不要填写 Topic 或协议前缀。
发布超时时间5000 毫秒单次发布等待上限,填写正整数。超时会交给目标缓存和失败策略处理。

安全认证

参数默认值如何配置
用户名仅 Broker 要求 SASL 认证时填写;明文匿名集群留空。
密码与 SASL 用户名配套;不要写入截图、模板或日志。
安全协议Plaintext选择与 Broker Listener 完全一致的 PlaintextSslSaslPlaintextSaslSsl
SASL机制Plain选择 Broker 配置的 GssapiPlainScramSha256ScramSha512OAuthBearer

安全协议、SASL 机制、用户名和密码必须成套配置;只填写用户名和密码不会自动启用 SASL。

消息配置

参数默认值如何配置
设备Topic模板设备消息 Topic。留空则不发布设备模型;填写固定 Topic 或 ${字段名} 模板。
变量Topic模板ThingsGateway/Variable变量消息 Topic。建议使用 ThingsGateway/Variable/${DeviceName} 按设备分组。
报警Topic模板报警消息 Topic。留空则不发布报警模型。
插件事件Topic模板插件事件 Topic。留空则不发布插件事件模型。

Topic 模板中的 ${字段名} 必须能在对应上传实体或实体脚本结果中找到;Topic 模板只决定路由,不改变转发组变量范围。

目标变量属性

Kafka 生产者提供 数据1数据10 十个可选文本字段,默认均为空。它们用于保存当前目标与变量的附加配置;Kafka 发送器不会自动把这些字段加入消息,只有你的实体脚本或后续自定义处理明确读取时才会产生作用。变量是否进入目标、别名、触发和分批仍在数据转发中配置。

参数默认值如何配置
数据1~数据10按项目约定填写文本;不需要时全部留空。不要把密码、Token 或其它敏感信息放入变量属性。

数据与脚本

参数默认值如何配置
详细日志关闭开启后记录每次上传的详细内容或计数,适合短时间排查;高频生产环境保持关闭。
JSON缩进格式化开启默认 JSON 带缩进,便于阅读;追求更小消息体时关闭。
JSON忽略Null开启开启后序列化时省略值为 null 的字段;需要保留空字段时关闭。
设备列表上传开启设备数据按列表合并发布;关闭后逐条发布。设备 Topic 为空时不会发布设备模型。
变量列表上传开启变量数据按列表合并发布;关闭后逐条发布。
变量字典上传关闭仅变量列表上传开启时生效;开启后按 DeviceName → Name → Value 组织字典。
报警列表上传开启报警数据按列表合并发布;关闭后逐条发布。
报警字典上传关闭仅报警列表上传开启时生效;开启后按设备和变量组织字典。
插件事件列表上传开启插件事件按列表合并发布;关闭后逐条发布。
设备实体脚本从表达式选择器选择脚本,将设备对象转换后再做 Topic 分组和消息序列化。
变量实体脚本选择脚本,将变量对象转换后再做 Topic 分组和消息序列化。
报警实体脚本选择脚本,将报警对象转换后再做 Topic 分组和消息序列化。
插件事件实体脚本选择脚本,将插件事件对象转换后再做 Topic 分组和消息序列化。

脚本输出必须仍能提供 Topic 模板和上传模板所引用的字段。脚本只改变发送对象,不改变转发组的变量范围。

上传模板配置

“上传模板配置”是一个嵌套对象,每类实体分别选择 TextJson 模式并填写模板正文;模板留空时使用默认 JSON 序列化。

配置项默认值如何配置
变量模板模式 / 变量内容模板Text / 空变量消息选择文本或 JSON,正文中使用 ${字段名}
设备模板模式 / 设备内容模板Text / 空设备消息选择文本或 JSON,正文中使用 ${字段名}
报警模板模式 / 报警内容模板Text / 空报警消息选择文本或 JSON,正文中使用 ${字段名}
插件事件模板模式 / 插件事件内容模板Text / 空插件事件选择文本或 JSON,正文中使用 ${字段名}

先在页面执行模板预览,再保存目标。JSON 模式必须生成合法 JSON;Text 模式不会替你补充引号或转义。

缓存与容量

参数默认值如何配置
启用失败重试缓存开启建议保持开启,发布失败的数据在恢复后自动补发;关闭后失败数据不再保留重试。
缓存文件最大行数262144CacheDB 出站队列上限,超过后删除最旧数据。按断网时长和消息速率规划磁盘。
上传分片大小2000每次从缓存取出的最大记录数。Kafka Broker 较慢时减小。
内存队列上限100000内存缓冲记录上限,超过后优先转入 CacheDB;持续超限仍可能丢弃旧数据。
过滤离线数据关闭开启后目标出队时过滤离线变量;转发组也有同名过滤,需同时检查。
并发上传数量1Kafka 生产者内部使用单锁保护发布,建议保持 1,不要仅为提高吞吐盲目增大。

模板可用字段

在“上传模板配置”中选择 Text 或 JSON 模式,使用 ${字段名} 插入对应实体或实体脚本输出的属性。保存前先执行模板预览。

数据类型可用字段
变量IdNameDeviceNameValueRawValueLastSetValueCollectGroupCollectTimeCreateTimeChangeTimeIsOnlineDataTypeUnitRegisterAddressOtherMethodDescriptionProtectTypeRpcWriteEnableRemark1Remark5ValueInitedIsMemory
设备IdNameActiveTimeDeviceStatusPluginNameDescriptionLastErrorMessageRemark1Remark5
报警AlarmIdVariableIdNameDeviceNameAlarmCodeAlarmLevelAlarmLimitAlarmTextRecoveryCodeAlarmTimeEventTimeFinishTimeConfirmTimeConfirmTextAlarmTypeEventTypeRemark1Remark5
插件事件DeviceNameObjectValue

实体脚本会先改变上传对象,再参与 Topic 分组和模板渲染;占位符必须与脚本输出一致。未开启对应 Topic 模板时,该类实体不会进入 Kafka 发布队列。

目标调试

进入“开发配置 → 数据转发”,选择 Kafka 目标,打开“调试 → 协议调试 · Kafka”。

功能作用
Kafka 协议调试使用测试 Topic 和小型消息检查连接、发布与返回结果。

Kafka 生产者专用调试面板

验证方法

  1. 确认 Bootstrap Server、安全协议和 SASL 设置。
  2. 在模板预览中检查测试 Topic 和小型 JSON 正文。
  3. 启用目标,让一个转发范围内的测试变量产生新值。
  4. 使用测试消费者核对 Topic、消息正文、分区、条数和时间。
  5. 检查目标日志中是否存在认证、超时或发布失败。

常见问题

现象检查方法
消费者没有消息服务地址、Topic 模板、集群 ACL、目标状态、消费者组和偏移。
认证失败安全协议、SASL 机制、用户名、密码和 Broker Listener。
发布超时Broker 网络、分区 Leader、集群负载和发布超时。
Topic 不正确Topic 模板、${字段名}、实体脚本输出和当前目标。
正文无效模板预览、JSON 语法、占位符和实体脚本输出。
变量属性没有进入消息数据1数据10 是预留元数据,Kafka 默认不会自动序列化;需要把值发出时应在实体脚本或上传模板中显式生成字段。
发布失败后数据消失检查“启用失败重试缓存”、CacheDB 待处理数量、内存队列上限和缓存文件最大行数。

相关操作

  • 数据转发:配置转发组、触发、缓存和公共目标操作。
  • 插件索引:查找其它采集或数据转发插件。