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.*
| クラス | 責務 |
|---|---|
KafkaDriver | ProtocolDriver 実装 — 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 protocol | payload | 使用場面 |
|---|---|---|---|
MQTT_CLIENT | MQTT v3.1.1 | 任意 (通常 JSON/平文) | IIoT broker (HiveMQ / Mosquitto / EMQX)、QoS 0/1/2 |
WEBSOCKET | ws/wss | 任意 (通常 JSON) | streaming server / proprietary push API |
KAFKA (本ページ) | Kafka wire protocol | record key/value (通常 JSON String) | high-throughput streaming、replay (offset)、consumer group |
SPARKPLUG_B | MQTT + Protobuf | Sparkplug B Payload | Cirrus 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|value | scalar (JSON でないメッセージの場合のみ) |
topic|<json-key> | top-level JSON key ごとの値 (object の場合のみ) |
PATH モードは cache を経由せず、rawByTopic の最後の raw メッセージを毎回 JSONPointer.queryFrom(...) で evaluate します。
オプション
| キー | デフォルト値 | 説明 |
|---|---|---|
group-id | plantpulse-edge-<opc_id> | consumer group id。offset commit の単位。 |
security-protocol | PLAINTEXT | PLAINTEXT / 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-offset | latest | earliest / 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 IAM —
aws-msk-iam-authjar の追加が必要です (現在は未同梱)。今後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 安全性
getDebugBaseUrl—kafka://、デフォルトポート (9092)- オプションのパース —
group-id/security-protocol/sasl-*/auto-offset、options=null 安全性 - 4-mode デコード — cache を直接埋めて KEY / PATH / RAW / SCALAR /
.valuesuffix をすべて検証 close()未接続時の安全性
test/java/plantpulse/driver/protocol/common/JsonExtractTest.java — 23 テスト (4-mode parser/cache/extract の統合検証)。
実 broker に依存する統合テストは別途です — 本単体テストは外部 Kafka なしで動作します。