メインコンテンツまでスキップ

plantpulse-messaging (メッセージング層)

役割

PlantPulse の神経系として機能するメッセージブローカー層です。産業現場から入ってくるセンサー データを収集し、モジュール間のイベントを伝えます。ブローカーは 2 つです — 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)

ブローカーが読む設定ファイルは 1 つだけ

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 リクエストを送って応答を得るか」を確認します。起動直後の DOWN (port 9092 is bound but … did not answer) は 30 秒後にもう一度確認してください。

チューニング ポイント

すべて /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 が意図的に除外されます)

よくある問題

症状原因対処
別ボックスのクライアントがブートストラップ後に切断されるアドバタイズ アドレスがコンテナ内部アドレスノードファイル PP_MASTER_IP設定を変更する方法
pd status messagingDOWN (advertises 127.0.0.1…)ノードファイルにアドレスがない同じ対処。レンダリング exit 8 も同じ原因
コンシューマー ラグが増加するコンシューマが遅い · 停止しているpd flowMEMB · 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)外部公開時に使用

関連ドキュメント