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

Apache Kafka ドライバ — 技術リファレンス

PlantPulse Edge の Apache Kafka ドライバは kafka-clients 4.1.2 ライブラリをラップした Java 実装です。 外部 streaming platform (Apache Kafka / Confluent Cloud / AWS MSK / Azure Event Hubs Kafka API) の topic を consumer として subscribe し、record を in-memory キャッシュに保持して、タグの polling cycle ごとにキャッシュ値を返します。 write は producer.send (async) です。

ソース: plantpulse.driver.protocol.kafka.*

クラス責務
KafkaDriverProtocolDriver 実装 — connect/read/write/close、topic キャッシュ、lazy subscribe、background poller thread

ライブラリ: lib/kafka-clients-4.1.2.jar


MQTT / WebSocket との違い

PlantPulse は 3 種類の message-based ドライバ を提供します — いずれも同一の JsonExtract 4-mode JSON デコーダを共有します。

ドライバwire protocolpayload使用場面
MQTT_CLIENTMQTT v3.1.1任意 (通常 JSON/平文)IIoT broker (HiveMQ / Mosquitto / EMQX)、QoS 0/1/2
WEBSOCKETws/wss任意 (通常 JSON)streaming server / proprietary push API
KAFKA (本ページ)Kafka wire protocolrecord key/value (通常 JSON String)high-throughput streaming、replay (offset)、consumer group
SPARKPLUG_BMQTT + ProtobufSparkplug B PayloadCirrus Link / Tahu / Ignition 標準 (MQTT 上の IIoT spec)

Kafka は persistent log + offset replay が核心です — 新規 consumer が auto-offset=earliest で登録すると 過去のデータから受信できます (broker の retention の範囲内)。


動作フロー

broker ──ConsumerRecord──> KafkaConsumer.poll() ──> pollerLoop (background thread)

├─ rawByTopic.put(topic, value)
└─ JsonExtract.updateCache(topicCache, topic, value)

PlantPulse 폴러 ──read(addr)──> JsonExtract.parse(addr.address) ──> JsonExtract.extract(...) (in-memory hit)
write(addr, value) ──KafkaProducer.send(ProducerRecord(topic, value))──> broker

read の初回呼び出しで cache miss の場合は lazy subscribe (consumer.subscribe(...) 呼び出し) 後に空文字列を返します — 次の cycle から実際の値が入ります。tag_map に登録された topic は connect 直後に preSubscribeFromTagMap() で一括 subscribe されるため、最初の cycle でただちに値が届く場合もあります。


4-mode JSON デコード spec

JsonExtract ユーティリティが MQTT / WebSocket / Kafka の 3 ドライバのメッセージデコードを統一します。

モードaddress 形式動作
SCALAR<topic> または <topic>.valueメッセージ全体を String として扱う。JSON object の場合は raw フォールバック。
KEY<topic>:<json-key>top-level JSON object の 1 つの key の値
PATH<topic>:$.<json-path>JSON Pointer を動的に evaluate (例: $.data.tags.T1/data/tags/T1)
RAW<topic>:_raw_最後の raw メッセージ (デバッグ用)

cache 構造 — 単一の Map<String, String> を使用し、key は topic + "|" + field:

key意味
topic|_raw_最後の raw メッセージ
topic|valuescalar (JSON でないメッセージの場合のみ)
topic|<json-key>top-level JSON key ごとの値 (object の場合のみ)

PATH モードは cache を経由せず、rawByTopic の最後の raw メッセージを毎回 JSONPointer.queryFrom(...) で evaluate します。


オプション

キーデフォルト値説明
group-idplantpulse-edge-<opc_id>consumer group id。offset commit の単位。
security-protocolPLAINTEXTPLAINTEXT / SSL / SASL_PLAINTEXT / SASL_SSL
sasl-mechanism(なし)PLAIN / SCRAM-SHA-256 / SCRAM-SHA-512
sasl-jaas-config(なし)JAAS 設定 (例: ...PlainLoginModule required username="u" password="p";)
auto-offsetlatestearliest / latest — 新規 group の開始位置

enable.auto.commit=true が常に設定され、broker が自動的に offset を commit します (デフォルト 5 秒周期)。


接続ライフサイクル

メソッド動作
connect()options 適用 → consumer / producer 生成 → poller thread 開始 → tag_map 一括 subscribe
pollerLoop()consumer.poll(500ms) 繰り返し → record ごとに JsonExtract.updateCache(...)
read()parse → cache miss なら lazy subscribe → cache hit なら値を返す
write()producer.send(new ProducerRecord<>(topic, value)) (async)
close()consumer.wakeup() → thread join → consumer/producer close → cache clear

制限事項および今後の作業

  • 再接続 — kafka-clients は独自の再接続ロジックを持ちますが、broker 全体が down して long-lived な poll() が失敗する場合、ドライバ側で setConnected(false) を個別にマーキングしません → 稼働状態の監視は broker side metric を推奨します。
  • AWS MSK IAMaws-msk-iam-auth jar の追加が必要です (現在は未同梱)。今後 aws-msk-iam-auth-1.x.jar を lib/ に配置し、sasl.client.callback.handler.class をオプションとして公開することを検討します。
  • TLS truststore — ドライバオプションでは指定できず、JVM のシステム truststore を使用します。自己署名のプライベート CA を使う broker の場合は、CA 証明書を JVM の truststore に追加するか、ゲートウェイ起動オプションに -Djavax.net.ssl.trustStore=… を追加してください。
  • partition / key 指定 producer — write は常に key なしで送信するため、パーティションは round-robin で決定されます。つまり 同一タグの書き込みが順序どおりに到達する保証はありません。 パーティション単位の順序が必要な場合は、topic をパーティション 1 個で作成するか、順序保証を consumer 側で処理してください。
  • transactional producer — 未対応。現在 idempotent producer も明示的に enable していません。

テストの場所

test/java/plantpulse/driver/protocol/kafka/KafkaDriverTest.java — 18 テスト。

  • クラスロード / 階層 (BaseProtocolDriver / ProtocolDriver)
  • 初期状態 / メタ (isConnected/isWriteSupported/isExternal/getDriverSource)
  • 未接続時の read/write/null 安全性
  • getDebugBaseUrlkafka://、デフォルトポート (9092)
  • オプションのパース — group-id / security-protocol / sasl-* / auto-offset、options=null 安全性
  • 4-mode デコード — cache を直接埋めて KEY / PATH / RAW / SCALAR / .value suffix をすべて検証
  • close() 未接続時の安全性

test/java/plantpulse/driver/protocol/common/JsonExtractTest.java — 23 テスト (4-mode parser/cache/extract の統合検証)。

実 broker に依存する統合テストは別途です — 本単体テストは外部 Kafka なしで動作します。