Apache Kafka Driver — Technische Referenz
Der Apache Kafka Treiber von PlantPulse Edge ist eine Java-Implementierung, die die Bibliothek kafka-clients 4.1.2 kapselt. Er abonniert als Consumer Topics einer externen Streaming Platform (Apache Kafka / Confluent Cloud / AWS MSK / Azure Event Hubs Kafka API), hält die Records in einem In-Memory-Cache und gibt bei jedem Polling-Zyklus eines Tags den Cache-Wert zurück. Schreibvorgänge erfolgen über producer.send (async).
Quelle: plantpulse.driver.protocol.kafka.*
| Klasse | Verantwortung |
|---|---|
KafkaDriver | Implementierung von ProtocolDriver — connect/read/write/close, Topic-Cache, Lazy Subscribe, Background-Poller-Thread |
Bibliothek: lib/kafka-clients-4.1.2.jar
Unterschiede zu MQTT / WebSocket
PlantPulse stellt drei message-basierte Treiber bereit — alle nutzen denselben 4-Mode-JSON-Decoder JsonExtract.
| Treiber | Wire Protocol | Payload | Einsatzfall |
|---|---|---|---|
MQTT_CLIENT | MQTT v3.1.1 | beliebig (meist JSON/Klartext) | IIoT-Broker (HiveMQ / Mosquitto / EMQX), QoS 0/1/2 |
WEBSOCKET | ws/wss | beliebig (meist JSON) | Streaming-Server / proprietäre Push-API |
KAFKA (diese Seite) | Kafka Wire Protocol | Record Key/Value (meist JSON-String) | High-Throughput-Streaming, Replay (Offset), Consumer Group |
SPARKPLUG_B | MQTT + Protobuf | Sparkplug B Payload | Cirrus Link / Tahu / Ignition Standard (IIoT-Spezifikation über MQTT) |
Der Kern von Kafka ist Persistent Log + Offset Replay — ein neuer Consumer, der sich mit auto-offset=earliest registriert,
kann auch zurückliegende Daten empfangen (innerhalb der Retention des Brokers).
Ablauf
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
Liegt beim ersten read-Aufruf ein Cache Miss vor, erfolgt ein Lazy Subscribe (Aufruf von consumer.subscribe(...)) und es wird ein leerer String zurückgegeben —
ab dem nächsten Zyklus kommen die tatsächlichen Werte an. Die in tag_map registrierten Topics werden unmittelbar nach connect
per preSubscribeFromTagMap() gesammelt abonniert, sodass bereits im ersten Zyklus Werte eintreffen können.
Spezifikation der 4-Mode-JSON-Dekodierung
Das Hilfsmodul JsonExtract vereinheitlicht die Nachrichtendekodierung der drei Treiber MQTT / WebSocket / Kafka.
| Modus | Adressformat | Verhalten |
|---|---|---|
| SCALAR | <topic> oder <topic>.value | Gesamte Nachricht als String. Bei einem JSON-Objekt Rückfall auf raw. |
| KEY | <topic>:<json-key> | Wert eines Keys im Top-Level-JSON-Objekt |
| PATH | <topic>:$.<json-path> | Dynamische Auswertung eines JSON Pointer (z. B. $.data.tags.T1 → /data/tags/T1) |
| RAW | <topic>:_raw_ | Letzte Rohnachricht (Debugging) |
Cache-Struktur — es wird eine einzige Map<String, String> verwendet, Key ist topic + "|" + field:
| Key | Bedeutung |
|---|---|
topic|_raw_ | Letzte Rohnachricht |
topic|value | Scalar (nur bei Nachrichten, die kein JSON sind) |
topic|<json-key> | Wert je Top-Level-JSON-Key (nur bei Objekten) |
Der PATH-Modus umgeht den Cache und wertet die letzte Rohnachricht aus rawByTopic jedes Mal mit JSONPointer.queryFrom(...) aus.
Optionen
| Key | Standardwert | Beschreibung |
|---|---|---|
group-id | plantpulse-edge-<opc_id> | Consumer-Group-ID. Einheit für den Offset-Commit. |
security-protocol | PLAINTEXT | PLAINTEXT / SSL / SASL_PLAINTEXT / SASL_SSL |
sasl-mechanism | (keiner) | PLAIN / SCRAM-SHA-256 / SCRAM-SHA-512 |
sasl-jaas-config | (keiner) | JAAS-Konfiguration (z. B. ...PlainLoginModule required username="u" password="p";) |
auto-offset | latest | earliest / latest — Startpunkt für neue Groups |
enable.auto.commit=true ist immer gesetzt, sodass der Broker Offsets automatisch committet (standardmäßig alle 5 Sekunden).
Verbindungslebenszyklus
| Methode | Verhalten |
|---|---|
connect() | Optionen anwenden → Consumer / Producer erzeugen → Poller-Thread starten → tag_map gesammelt abonnieren |
pollerLoop() | consumer.poll(500ms) in Schleife → je Record JsonExtract.updateCache(...) |
read() | Parsen → bei Cache Miss Lazy Subscribe → bei Cache Hit Wert zurückgeben |
write() | producer.send(new ProducerRecord<>(topic, value)) (async) |
close() | consumer.wakeup() → Thread-Join → Consumer/Producer schließen → Cache leeren |
Einschränkungen und offene Punkte
- Wiederverbindung — kafka-clients besitzt eine eigene Reconnect-Logik; fällt jedoch der gesamte Broker aus und schlägt ein langlebiger
poll()fehl, markiert der TreibersetConnected(false)nicht gesondert → für das Betriebsmonitoring werden broker-seitige Metriken empfohlen. - AWS MSK IAM — erfordert das zusätzliche jar
aws-msk-iam-auth(derzeit nicht enthalten). Künftig wird geprüft,aws-msk-iam-auth-1.x.jarin lib/ abzulegen undsasl.client.callback.handler.classals Option verfügbar zu machen. - TLS-Truststore — nicht über Treiberoptionen konfigurierbar; es wird der System-Truststore der JVM verwendet. Bei Brokern mit selbstsignierter, privater CA das CA-Zertifikat in den JVM-Truststore aufnehmen oder den Startoptionen des Gateway
-Djavax.net.ssl.trustStore=…hinzufügen. - Producer mit Partition-/Key-Angabe — write sendet stets ohne Key, die Partition wird daher per Round-Robin bestimmt. Das heißt: es ist nicht garantiert, dass Schreibvorgänge desselben Tags in der richtigen Reihenfolge ankommen. Wird eine Reihenfolge je Partition benötigt, das Topic mit nur einer Partition anlegen oder die Reihenfolge consumer-seitig sicherstellen.
- Transaktionaler Producer — nicht unterstützt. Derzeit wird auch der idempotente Producer nicht explizit aktiviert.
Speicherort der Tests
test/java/plantpulse/driver/protocol/kafka/KafkaDriverTest.java — 18 Tests.
- Klassenladen / Hierarchie (
BaseProtocolDriver/ProtocolDriver) - Initialzustand / Meta (
isConnected/isWriteSupported/isExternal/getDriverSource) - read/write ohne Verbindung / Null-Sicherheit
getDebugBaseUrl—kafka://, Standardport (9092)- Optionsparsing —
group-id/security-protocol/sasl-*/auto-offset, Sicherheit bei options=null - 4-Mode-Dekodierung — Cache direkt befüllt, Prüfung von KEY / PATH / RAW / SCALAR und dem Suffix
.value close()Sicherheit ohne Verbindung
test/java/plantpulse/driver/protocol/common/JsonExtractTest.java — 23 Tests (Integrationsprüfung von 4-Mode-Parser/Cache/Extract).
Integrationstests mit einem realen Broker sind separat — diese Unit-Tests laufen ohne externes Kafka.