Kafka 生产者
插件用途
将转发组产生的变量、设备、报警和插件事件发布到 Apache Kafka Topic。
转发组范围、触发、分批、缓存和启停见数据转发。本页只说明 Kafka 连接、安全认证、Topic 模板、上传模板和专用调试。
功能入口
进入“开发配置 → 数据转发”,按以下顺序配置:
- 新增或编辑转发组,先保存变量范围、触发模式、定时间隔、在线过滤和批处理策略。
- 新增目标,选择“Kafka 生产者”,填写目标名称、启用状态、日志级别和启动超时。
- 打开“目标属性”,依次填写 Kafka 地址、安全认证、消息配置、数据与脚本以及缓存容量。
- 保存并启用转发组和目标,再从“目标调试”验证发布结果。
目标基本信息
| 参数 | 默认值 | 如何配置 |
|---|---|---|
| 所属转发组 | - | 必须选择已保存的转发组;目标只接收该组范围内的变量。 |
| 目标名称 | - | 必填,同一转发组内不能重复。建议包含环境和 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模板 | 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 模板和上传模板所引用的字段。脚本只改变发送对象,不改变转发组的变量范围。
上传模板配置
“上传模板配置”是一个嵌套对象,每类实体分别选择 Text 或 Json 模式并填写模板正文;模板留空时使用默认 JSON 序列化。
| 配置项 | 默认值 | 如何配置 |
|---|---|---|
| 变量模板模式 / 变量内容模板 | Text / 空 | 变量消息选择文本或 JSON,正文中使用 ${字段名}。 |
| 设备模板模式 / 设备内容模板 | Text / 空 | 设备消息选择文本或 JSON,正文中使用 ${字段名}。 |
| 报警模板模式 / 报警内容模板 | Text / 空 | 报警消息选择文本或 JSON,正文中使用 ${字段名}。 |
| 插件事件模板模式 / 插件事件内容模板 | Text / 空 | 插件事件选择文本或 JSON,正文中使用 ${字段名}。 |
先在页面执行模板预览,再保存目标。JSON 模式必须生成合法 JSON;Text 模式不会替你补充引号或转义。
缓存与容量
| 参数 | 默认值 | 如何配置 |
|---|---|---|
| 启用失败重试缓存 | 开启 | 建议保持开启,发布失败的数据在恢复后自动补发;关闭后失败数据不再保留重试。 |
| 缓存文件最大行数 | 262144 | CacheDB 出站队列上限,超过后删除最旧数据。按断网时长和消息速率规划磁盘。 |
| 上传分片大小 | 2000 | 每次从缓存取出的最大记录数。Kafka Broker 较慢时减小。 |
| 内存队列上限 | 100000 | 内存缓冲记录上限,超过后优先转入 CacheDB;持续超限仍可能丢弃旧数据。 |
| 过滤离线数据 | 关闭 | 开启后目标出队时过滤离线变量;转发组也有同名过滤,需同时检查。 |
| 并发上传数量 | 1 | Kafka 生产者内部使用单锁保护发布,建议保持 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 和小型消息检查连接、发布与返回结果。 |

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