跳到主要内容

Kafka 生产者

插件用途​

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

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

功能入口​

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

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

目标基本信息​

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

目标属性​

连接配置​

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

安全认证​

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

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

消息配置​

参数默认值如何配置
设备Topic模板空设备消息 Topic。留空则不发布设备模型;填写固定 Topic 或 ${字段名} 模板。
变量Topic模板TGateway/Variable变量消息 Topic。建议使用 TGateway/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 模板和上传模板所引用的字段。脚本只改变发送对象,不改变转发组的变量范围。

上传模板配置​

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

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

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

缓存与容量​

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

模板可用字段​

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

数据类型可用字段
变量Id、Name、DeviceName、Value、RawValue、LastSetValue、CollectGroup、CollectTime、CreateTime、ChangeTime、IsOnline、DataType、Unit、RegisterAddress、OtherMethod、Description、ProtectType、RpcWriteEnable、Remark1~Remark5、ValueInited、IsMemory
设备Id、Name、ActiveTime、DeviceStatus、PluginName、Description、LastErrorMessage、Remark1~Remark5
报警AlarmId、VariableId、Name、DeviceName、AlarmCode、AlarmLevel、AlarmLimit、AlarmText、RecoveryCode、AlarmTime、EventTime、FinishTime、ConfirmTime、ConfirmText、AlarmType、EventType、Remark1~Remark5
插件事件DeviceName、ObjectValue

实体脚本会先改变上传对象,再参与 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 待处理数量、内存队列上限和缓存文件最大行数。

相关操作​

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