Zum Hauptinhalt springen

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.*

KlasseVerantwortung
KafkaDriverImplementierung 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.

TreiberWire ProtocolPayloadEinsatzfall
MQTT_CLIENTMQTT v3.1.1beliebig (meist JSON/Klartext)IIoT-Broker (HiveMQ / Mosquitto / EMQX), QoS 0/1/2
WEBSOCKETws/wssbeliebig (meist JSON)Streaming-Server / proprietäre Push-API
KAFKA (diese Seite)Kafka Wire ProtocolRecord Key/Value (meist JSON-String)High-Throughput-Streaming, Replay (Offset), Consumer Group
SPARKPLUG_BMQTT + ProtobufSparkplug B PayloadCirrus 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.

ModusAdressformatVerhalten
SCALAR<topic> oder <topic>.valueGesamte 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:

KeyBedeutung
topic|_raw_Letzte Rohnachricht
topic|valueScalar (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

KeyStandardwertBeschreibung
group-idplantpulse-edge-<opc_id>Consumer-Group-ID. Einheit für den Offset-Commit.
security-protocolPLAINTEXTPLAINTEXT / 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-offsetlatestearliest / 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

MethodeVerhalten
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 Treiber setConnected(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.jar in lib/ abzulegen und sasl.client.callback.handler.class als 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
  • getDebugBaseUrlkafka://, 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.