Apache Kafka
网关以客户端身份向外发起 outbound 连接,接入外部数据 streaming platform(Apache Kafka / Confluent Cloud / AWS MSK / Azure Event Hubs Kafka API)的 topic,进行 consumer subscribe / producer send。 收到的 record 会存入内存缓存,标签在每个采集周期读取缓存值。
| 场景 | 应使用哪种模式 |
|---|---|
| 外部 IT 系统通过 REST POST 发送数值时 | HTTP 推送 |
| 外部服务器通过 WebSocket push 时 | WebSocket Client |
| 外部 IIoT MQTT broker 通过 topic push 时 | MQTT Client |
| 以 consumer 接收 Kafka topic 时 | Apache Kafka(本页) |
| 直接读取厂内 PLC 数值时 | Modbus / OPC-UA 等 |
注册表单输入项
| 输入项 | 填写内容 | 示例 |
|---|---|---|
| IP 地址 | Kafka broker 主机(bootstrap) | kafka.example.com、10.0.0.50 |
| 端口 | broker 端口(明文 9092,TLS 9093) | 9092、9093 |
| GROUP ID | consumer group id(可选) | plant-floor-A、edge-line1 |
| SECURITY | 安全协议 | PLAINTEXT / SSL / SASL_PLAINTEXT / SASL_SSL |
| SASL MECH | SASL 机制(可选) | PLAIN / SCRAM-SHA-256 / SCRAM-SHA-512 |
| SASL JAAS | SASL 认证 JAAS 配置(可选) | org.apache.kafka.common.security.plain.PlainLoginModule required username="u" password="p"; |
| AUTO OFFSET | 新建 group 的起始位置 | latest(默认) / earliest |
| 采集周期 | 网关读取缓存值的周期(ms) | 1000 |
实际 broker URL:<host>:<port>(例如:kafka.example.com:9092)。bootstrap 服务器即使只注册一个主机,
broker 也会自动返回 cluster 信息。
未指定 GROUP ID 时,将按 plantpulse-edge-<opc_id> 格式自动生成。
标签的 PLC 地址表示法 —— 4-mode JSON
标签的 PLC 地址 = Kafka topic + 4-mode 解码器。与 MQTT / WebSocket 规格相同。
| 模式 | 格式 | 行为 |
|---|---|---|
| SCALAR | sensors/temp 或 sensors/temp.value | 将整条消息作为 String 处理。若消息为 JSON object,则回退为 raw。 |
| KEY | sensors/temp:temperature | top-level JSON key 的值(例如:{"temperature":25.3,"humidity":60} → 25.3) |
| PATH | sensors/temp:$.data.tags.T1 | JSON Pointer 动态求值(支持嵌套 key) |
| RAW | sensors/temp:_raw_ | 最后一条 record value 全文(调试用) |
read 的首次调用为 lazy subscribe(返回空字符串)。从下一个 polling cycle 起才会取到缓存值。
常见用例
| 用例 | 如何配置 |
|---|---|
| 厂内 Kafka cluster | host = bootstrap broker IP,9092(明文) |
| Confluent Cloud | host = pkc-XXX.region.aws.confluent.cloud,9092 + SASL_SSL + PLAIN mechanism + API key/secret JAAS |
| AWS MSK(IAM auth) | host = b-1.<cluster>...amazonaws.com,9098,IAM auth(需另行配置 JAAS) |
| Azure Event Hubs(Kafka API) | host = <ns>.servicebus.windows.net,9093 + SASL_SSL + PLAIN + Connection String JAAS |
| 新建 consumer group + earliest | 为接收全部历史数据,选择 AUTO OFFSET=earliest |
| 分离管理 group offset | 使用不同的 GROUP ID 注册 —— 同一 topic 也会按 group 分别管理 offset |
write(producer.send)
在标签页面或通过 REST API 写入数值时,producer 会将其作为 record 发送到对应 topic(async,fire-and-forget)。 value 的 String 将作为 record value 发送 —— JSON / 明文均可。
plc_address = factory/line1/cmd
value = ON
→ Kafka producer.send: topic="factory/line1/cmd" value="ON"
plc_address = factory/line1/cmd:temperature
value = 25.3
→ Kafka producer.send: topic="factory/line1/cmd" value="25.3" (콜론 뒤는 read decoder 정보, write 는 topic 만 사용)
常见问题与解决方法
| 现象 | 原因 | 解决方法 |
|---|---|---|
| 收不到数值 | broker 未发送 record | 用 kafka-console-consumer.sh --bootstrap-server <host>:9092 --topic <t> --from-beginning 直接接收并确认 |
[KAFKA] connect 실패: TimeoutException | 无法到达 bootstrap 服务器 | 验证 telnet <host> 9092。确认 advertised.listeners 是否为外部 IP(broker 侧) |
[KAFKA] poll error: Authentication failed | SASL 认证失败 | 重新确认 sasl-mechanism / sasl-jaas-config。password 中的特殊字符需转义 |
[KAFKA] poll error: SslHandshake | TLS truststore 缺失 / 过期 | 使用可信 CA 签发 broker 证书,或将其加入系统 truststore |
| 新建 group 收不到数据 | 设为 auto-offset=latest 但没有新 record | 临时改为 earliest,或等待发生 publish |