跳到主要内容

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 topicApache Kafka(本页)
直接读取厂内 PLC 数值时Modbus / OPC-UA

注册表单输入项

输入项填写内容示例
IP 地址Kafka broker 主机(bootstrap)kafka.example.com10.0.0.50
端口broker 端口(明文 9092,TLS 9093)90929093
GROUP IDconsumer group id(可选)plant-floor-Aedge-line1
SECURITY安全协议PLAINTEXT / SSL / SASL_PLAINTEXT / SASL_SSL
SASL MECHSASL 机制(可选)PLAIN / SCRAM-SHA-256 / SCRAM-SHA-512
SASL JAASSASL 认证 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 规格相同。

模式格式行为
SCALARsensors/tempsensors/temp.value将整条消息作为 String 处理。若消息为 JSON object,则回退为 raw。
KEYsensors/temp:temperaturetop-level JSON key 的值(例如:{"temperature":25.3,"humidity":60}25.3
PATHsensors/temp:$.data.tags.T1JSON Pointer 动态求值(支持嵌套 key)
RAWsensors/temp:_raw_最后一条 record value 全文(调试用)

read 的首次调用为 lazy subscribe(返回空字符串)。从下一个 polling cycle 起才会取到缓存值。


常见用例

用例如何配置
厂内 Kafka clusterhost = bootstrap broker IP,9092(明文)
Confluent Cloudhost = 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 未发送 recordkafka-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 failedSASL 认证失败重新确认 sasl-mechanism / sasl-jaas-config。password 中的特殊字符需转义
[KAFKA] poll error: SslHandshakeTLS truststore 缺失 / 过期使用可信 CA 签发 broker 证书,或将其加入系统 truststore
新建 group 收不到数据设为 auto-offset=latest 但没有新 record临时改为 earliest,或等待发生 publish

更详细的技术文档