plantpulse-messaging(消息层)
角色
PlantPulse 的神经系统——消息代理层。采集来自工业现场的传感器数据,转发模块间的事件。有两个代理:Kafka 和 HiveMQ(MQTT)。
| 项目 | 值 |
|---|---|
| 模块名 | plantpulse-messaging |
| 安装路径 | /opt/kopens/plantpulse-platform/plantpulse-messaging/ |
| 数据位置 | /data1/pp-data/kafka/kraft_combined_logs |
pd 服务 | messaging — 仅在 MASTER 节点 |
| 账户 | PP_MQ_USER / PP_MQ_PASSWORD — Kafka 和 HiveMQ 共享同一个值 |
浏览器实时推送改为 SSE 后,镜像中已移除 STOMP 代理。端口 61000 / 61004 和 PP_STOMP_* 不再提供服务。如旧防火墙规则中仍有残留,请清理。
架构
Kafka
KRaft 模式(无 ZooKeeper)的单代理兼控制器。
| 项目 | 值 |
|---|---|
| 端口 | 9092 (SASL_PLAINTEXT) · 9093 (controller,内部) · 9094 (SASL_SSL) |
| 认证 | SASL/PLAIN — 静态 JAAS(config/jaas.conf)。账户 PP_MQ_USER |
| 广播地址 | PP_KAFKA_ADVERTISED_HOST — 启动时由主机地址驱动 |
| 主题前缀 | PP_TOPIC_PREFIX (默认 pp) |
代理读取的文件只有一个
plantpulse-messaging/kafka/config/kafka.properties 是代理配置,从主机的 /etc/kopens/conf/kafka.properties.template 渲染。同目录中的 broker.properties · server.properties · controller.properties 是 Apache Kafka 发行版自带的示例文件,无人使用。不要对容器内的 advertised.listeners=PLAINTEXT://localhost:9092 等值感到惊讶。
模板的关键行(占位符在渲染时填充):
process.roles=broker,controller
listeners=SASL_PLAINTEXT://0.0.0.0:${PP_KAFKA_PORT},CONTROLLER://0.0.0.0:${PP_KAFKA_CONTROLLER_PORT:9093},SASL_SSL://0.0.0.0:${PP_KAFKA_TLS_PORT}
advertised.listeners=SASL_PLAINTEXT://${PP_KAFKA_ADVERTISED_HOST}:${PP_KAFKA_PORT},…
sasl.enabled.mechanisms=PLAIN
allow.everyone.if.no.acl.found=false
log.dirs=${PP_DATA_DIR}/kafka/kraft_combined_logs
log.retention.ms=3600000 # 1시간 — 버퍼이지 보관소가 아닙니다
log.segment.bytes=10485760 # 10MB
num.partitions=4
default.replication.factor=1 # 단일 브로커
message.max.bytes=157286400
auto.create.topics.enable=true
代理默认保留期为 1 小时。数据的永久存储是 Cassandra · 时序引擎 · Iceberg,Kafka 是其间的缓冲。每个主题的保留和«谁设置的这个值»由 pd retention 显示。
关键主题
使用前缀 pp(PP_TOPIC_PREFIX)。下表是平台代码使用的典型主题。
| 主题 | 用途 |
|---|---|
event · event-row | 服务器 ↔ 引擎事件总线 |
pp-tag-point | 标签值流 — 需为每个标签启用«Kafka 发布»才会填充 |
pp-tag-alarm · pp-asset-alarm · pp-asset-event | 报警 · 资产事件 |
pp-asset-data · pp-asset-aggregation · pp-asset-health-status | 资产数据和聚合 |
pp-batch | 批处理作业 |
pp-edge-status · pp-opc-status · pp-diagnostic | Edge · OPC · 诊断状态 |
pp-domain-changed-event | 元数据变更 |
pp-production-oee · pp-production-ram · pp-production-ems | 生产指标 |
运维命令
docker exec plantpulse-datalake pd status messaging # UP (advertised <address>:9092 answered ApiVersions)
docker exec plantpulse-datalake pd flow # 컨슈머 그룹별 lag
docker exec plantpulse-datalake pd node topic # 토픽 "event" describe (비밀번호 자동 주입)
docker exec plantpulse-datalake pd retention # 토픽별 보존과 세그먼트 크기
docker exec plantpulse-datalake pd restart messaging
要直接使用 Kafka 发行版的工具,需通过客户端配置文件传递 SASL 凭证。pd flow · pd node topic 可代您完成,请先使用它们。
代理会在开始处理请求前约 10 秒打开 9092。因此 pd 的健康检查不是看端口,而是«通过广播地址发送 ApiVersions 请求是否收到应答»。启动直后,请在 30 秒后重新检查 DOWN (port 9092 is bound but … did not answer)。
调优要点
全部在 /etc/kopens/conf/kafka.properties.template 中修改,用 restart-datalake.sh 生效。
| 项目 | 默认 | 备注 |
|---|---|---|
log.retention.ms | 1 小时 | 消费者可能长时间停止的场景下增加。按此占用磁盘空间 |
num.partitions | 4 | 新主题的分区数。与消费者线程数保持一致 |
num.io.threads · num.network.threads | 8 · 3 | NVMe 则将 io 增至 16 |
message.max.bytes | 150MB | 为大消息(文件队列)设置了较大值 |
HiveMQ (MQTT)
轻量级 IoT 设备、传感器和 Edge 代理连接的标准 MQTT 代理(Community Edition)。
| 项目 | 值 |
|---|---|
| 端口 | 1883 (明文) · 1884 (TLS) — 代理容器在主机上发布并转发至数据湖 |
| 认证 | username / password — 安全扩展读取 conf/auth.properties |
| 账户 | PP_MQ_USER / PP_MQ_PASSWORD (与 Kafka 相同值) |
| 配置 | mqtt/conf/config.xml ← /etc/kopens/conf/hivemq.xml.template |
TLS 监听器使用平台公共密钥库(/var/security/plantpulse/master/master.keystore.jks),不要求客户端证书(NONE)。
# 외부에서 발행 · 구독 (디버깅)
mosquitto_pub -h <server-ip> -p 1883 -u mq -P '<PP_MQ_PASSWORD>' -t test/topic -m "hello"
mosquitto_sub -h <server-ip> -p 1883 -u mq -P '<PP_MQ_PASSWORD>' -t 'test/#'
日志
docker exec plantpulse-datalake pd logs --lines 100 messaging
| 代理 | 文件 |
|---|---|
| Kafka | plantpulse-messaging/kafka/logs/ |
| HiveMQ | plantpulse-messaging/mqtt/logs/hivemq.log (event.log 是按消息审计日志,pd logs 特意省略) |
常见问题
| 症状 | 原因 | 解决方案 |
|---|---|---|
| 其他主机上的客户端在 bootstrap 后断开 | 广播地址是容器内部地址 | 节点文件 PP_MASTER_IP → 修改配置的方法 |
pd status messaging 是 DOWN (advertises 127.0.0.1…) | 节点文件中无地址 | 同上。与渲染 exit 8 原因相同 |
| 消费者 lag 增加 | 消费者慢 · 停止 | pd flow 中的 MEMB · LAG,消费者容器日志 |
| 磁盘满 | 增加了保留但消费者停止 | pd retention 中的 kafka 表,log.retention.ms |
| MQTT TLS handshake 失败 | 证书过期 · SAN 缺失 | 安全设置 — plantpulse-certs 会重新签发 |
| 轮换后仅 Kafka 认证失败 | PP_MQ_PASSWORD 是两个代理共用 | 用 passwd.sh PP_MQ_PASSWORD 同步两个 → 密码 |
安全 / 外部暴露
| 端口 | 建议 |
|---|---|
| 9092 (Kafka SASL 明文) | 仅限专网 |
| 9094 (Kafka SASL_SSL) | 外部暴露时使用 |
| 1883 (MQTT 明文) | 仅限专网 |
| 1884 (MQTT TLS) | 外部暴露时使用 |