Apache Kafka
外部データストリーミングプラットフォーム (Apache Kafka / Confluent Cloud / AWS MSK / Azure Event Hubs Kafka API) の topic に対して、ゲートウェイがクライアントとして outbound 接続し、consumer subscribe / producer send を行います。 到着した record は in-memory キャッシュに保存され、タグの収集周期ごとにキャッシュ値を読み取ります。
| 状況 | どのモードを使うか |
|---|---|
| 外部 IT システムが REST POST で値を送ってくる場合 | HTTP プッシュ |
| 外部サーバーが WebSocket で push する場合 | WebSocket Client |
| 外部 IIoT MQTT broker が topic に push する場合 | MQTT Client |
| Kafka topic を consumer として受信する場合 | Apache Kafka (このページ) |
| 社内 PLC の値を直接読み取る場合 | Modbus / OPC-UA など |
登録フォームの入力値
| 入力欄 | 何を入力するか | 例 |
|---|---|---|
| IP アドレス | Kafka broker ホスト (bootstrap) | kafka.example.com, 10.0.0.50 |
| ポート | broker ポート (平文 9092、TLS 9093) | 9092, 9093 |
| GROUP ID | consumer group id (オプション) | plant-floor-A, edge-line1 |
| SECURITY | セキュリティプロトコル | PLAINTEXT / SSL / SASL_PLAINTEXT / SASL_SSL |
| SASL MECH | SASL メカニズム (オプション) | PLAIN / SCRAM-SHA-256 / SCRAM-SHA-512 |
| SASL JAAS | SASL 認証 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 サーバーは 1 ホストだけ登録しても、
broker が cluster 情報を自動的に通知します。
GROUP ID を指定しない場合は plantpulse-edge-<opc_id> 形式で自動生成されます。
タグの PLC アドレス表記 — 4-mode JSON
タグの PLC アドレス = Kafka topic + 4-mode デコーダ。MQTT / WebSocket と同一の spec です。
| モード | 形式 | 動作 |
|---|---|---|
| SCALAR | sensors/temp または sensors/temp.value | メッセージ全体を String として扱う。メッセージが JSON object の場合は raw フォールバック。 |
| KEY | sensors/temp:temperature | top-level JSON key の値 (例: {"temperature":25.3,"humidity":60} → 25.3) |
| PATH | sensors/temp:$.data.tags.T1 | JSON Pointer による動的 evaluate (ネストされた key に対応) |
| RAW | sensors/temp:_raw_ | 最後の record value 全体 (デバッグ用) |
read の初回呼び出しは lazy subscribe となります (空文字列を返します)。次の polling cycle からキャッシュ値が入ります。
よく使われるケース
| ケース | 方法 |
|---|---|
| 社内 Kafka cluster | host = bootstrap broker IP、9092 (平文) |
| Confluent Cloud | host = 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 が該当 topic へ record として送信します (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 が record を送信していない | kafka-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 failed | SASL 認証失敗 | sasl-mechanism / sasl-jaas-config を再確認。password の特殊文字を escape |
[KAFKA] poll error: SslHandshake | TLS truststore の欠落 / 期限切れ | broker 証明書を信頼できる CA で発行、またはシステム truststore に追加 |
| 新規 group がデータを受信できない | auto-offset=latest だが新規 record が存在しない | 一時的に earliest へ変更、または publish が発生するまで待機 |