跳到主要内容

plantpulse-messaging(消息层)

角色

PlantPulse 的神经系统——消息代理层。采集来自工业现场的传感器数据,转发模块间的事件。有两个代理:KafkaHiveMQ(MQTT)

项目
模块名plantpulse-messaging
安装路径/opt/kopens/plantpulse-platform/plantpulse-messaging/
数据位置/data1/pp-data/kafka/kraft_combined_logs
pd 服务messaging — 仅在 MASTER 节点
账户PP_MQ_USER / PP_MQ_PASSWORDKafka 和 HiveMQ 共享同一个值
STOMP(ActiveMQ)已停用

浏览器实时推送改为 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.propertiesApache 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
Kafka 是 1 小时缓冲区

代理默认保留期为 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-diagnosticEdge · 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 可代您完成,请先使用它们。

9092 开放不等于已就绪

代理会在开始处理请求前约 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.ms1 小时消费者可能长时间停止的场景下增加。按此占用磁盘空间
num.partitions4新主题的分区数。与消费者线程数保持一致
num.io.threads · num.network.threads8 · 3NVMe 则将 io 增至 16
message.max.bytes150MB为大消息(文件队列)设置了较大值

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
代理文件
Kafkaplantpulse-messaging/kafka/logs/
HiveMQplantpulse-messaging/mqtt/logs/hivemq.log (event.log 是按消息审计日志,pd logs 特意省略)

常见问题

症状原因解决方案
其他主机上的客户端在 bootstrap 后断开广播地址是容器内部地址节点文件 PP_MASTER_IP修改配置的方法
pd status messagingDOWN (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)外部暴露时使用

相关文档