plantpulse-messaging (メッセージング層)
役割
PlantPulse の神経系として機能するメッセージブローカー層です。産業現場から入ってくるセンサー データを収集し、モジュール間のイベントを伝えます。ブローカーは 2 つです — 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) |
ブローカーが読む設定ファイルは 1 つだけ
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 リクエストを送って応答を得るか」を確認します。起動直後の DOWN (port 9092 is bound but … did not answer) は 30 秒後にもう一度確認してください。
チューニング ポイント
すべて /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 が意図的に除外されます) |
よくある問題
| 症状 | 原因 | 対処 |
|---|---|---|
| 別ボックスのクライアントがブートストラップ後に切断される | アドバタイズ アドレスがコンテナ内部アドレス | ノードファイル PP_MASTER_IP → 設定を変更する方法 |
pd status messaging が DOWN (advertises 127.0.0.1…) | ノードファイルにアドレスがない | 同じ対処。レンダリング exit 8 も同じ原因 |
| コンシューマー ラグが増加する | コンシューマが遅い · 停止している | pd flow の MEMB · LAG、コンシューマ コンテナログ |
| ディスク満杯 | 保持期間を延長したがコンシューマが停止している | pd retention の kafka テーブル、log.retention.ms |
| MQTT TLS ハンドシェイク失敗 | 証明書期限切れ · SAN 不足 | セキュリティ設定 — plantpulse-certs が再発行 |
| ローテーション後 Kafka のみ認証失敗 | PP_MQ_PASSWORD は 2 つのブローカーで共有 | passwd.sh PP_MQ_PASSWORD で両方を一緒に → パスワード |
セキュリティ / 外部公開
| ポート | 推奨 |
|---|---|
| 9092 (Kafka SASL 平文) | プライベートネットワークのみ |
| 9094 (Kafka SASL_SSL) | 外部公開時に使用 |
| 1883 (MQTT 平文) | プライベートネットワークのみ |
| 1884 (MQTT TLS) | 外部公開時に使用 |