Zum Hauptinhalt springen

Apache Kafka

Das Gateway verbindet sich als Client per Outbound-Verbindung mit den Topics einer externen Daten-Streaming-Plattform (Apache Kafka / Confluent Cloud / AWS MSK / Azure Event Hubs Kafka API) und führt Consumer-Subscribe / Producer-Send aus. Eingehende Records werden in einem In-Memory-Cache abgelegt; bei jedem Erfassungsintervall des Tags wird der Cache-Wert ausgelesen.

SituationWelcher Modus ist zu verwenden
Ein externes IT-System sendet Werte per REST POSTHTTP-Push
Ein externer Server pusht per WebSocketWebSocket Client
Ein externer IIoT-MQTT broker pusht auf ein TopicMQTT Client
Kafka topic wird als Consumer empfangenApache Kafka (diese Seite)
Werte werksinterner PLCs werden direkt gelesenModbus / OPC-UA usw.

Eingabefelder des Registrierungsformulars

EingabefeldWas wird eingetragenBeispiel
IP-AdresseKafka-broker-Host (bootstrap)kafka.example.com, 10.0.0.50
Portbroker-Port (unverschlüsselt 9092, TLS 9093)9092, 9093
GROUP IDconsumer group id (optional)plant-floor-A, edge-line1
SECURITYSicherheitsprotokollPLAINTEXT / SSL / SASL_PLAINTEXT / SASL_SSL
SASL MECHSASL-Mechanismus (optional)PLAIN / SCRAM-SHA-256 / SCRAM-SHA-512
SASL JAASJAAS-Konfiguration für die SASL-Authentifizierung (optional)org.apache.kafka.common.security.plain.PlainLoginModule required username="u" password="p";
AUTO OFFSETStartpunkt für neue Groupslatest (Standard) / earliest
ErfassungsintervallIntervall, in dem das Gateway die zwischengespeicherten Werte ausliest (ms)1000

Tatsächliche broker-URL: <host>:<port> (z. B. kafka.example.com:9092). Es genügt, einen einzigen Host als bootstrap-Server einzutragen — der broker übermittelt die Cluster-Informationen automatisch.

Ohne Angabe einer GROUP ID wird sie automatisch im Format plantpulse-edge-<opc_id> erzeugt.


PLC-Adressnotation des Tags — 4-Mode-JSON

PLC-Adresse des Tags = Kafka topic + 4-Mode-Decoder. Identische Spezifikation wie bei MQTT / WebSocket.

ModusFormatVerhalten
SCALARsensors/temp oder sensors/temp.valueGesamte Nachricht als String. Ist die Nachricht ein JSON-Objekt, erfolgt ein Raw-Fallback.
KEYsensors/temp:temperatureWert eines Top-Level-JSON-Keys (z. B. {"temperature":25.3,"humidity":60}25.3)
PATHsensors/temp:$.data.tags.T1Dynamische Auswertung per JSON Pointer (verschachtelte Keys unterstützt)
RAWsensors/temp:_raw_Vollständiger Value des letzten Records (Debugging)

Der erste read-Aufruf löst ein Lazy Subscribe aus (liefert einen leeren String). Ab dem nächsten Polling-Zyklus stehen die Cache-Werte zur Verfügung.


Häufige Anwendungsfälle

AnwendungsfallVorgehen
Werksinternes Kafka-Clusterhost = bootstrap-broker-IP, 9092 (unverschlüsselt)
Confluent Cloudhost = pkc-XXX.region.aws.confluent.cloud, 9092 + SASL_SSL + Mechanismus PLAIN + API-Key/Secret-JAAS
AWS MSK (IAM auth)host = b-1.<cluster>...amazonaws.com, 9098, IAM auth (separate JAAS-Konfiguration erforderlich)
Azure Event Hubs (Kafka API)host = <ns>.servicebus.windows.net, 9093 + SASL_SSL + PLAIN + Connection-String-JAAS
Neue consumer group + earliestAUTO OFFSET=earliest wählen, um alle vorhandenen Daten zu empfangen
Getrennte Verwaltung der Group-OffsetsMit anderer GROUP ID registrieren — auch beim selben Topic wird der Offset je Group separat verwaltet

write (producer.send)

Wird ein Wert über die Tag-Seite oder die REST API geschrieben, sendet der Producer ihn als Record an das entsprechende Topic (async, fire-and-forget). Der String des Werts wird als Record-Value übertragen — JSON und Klartext sind gleichermaßen möglich.

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 만 사용)

Häufige Probleme und Lösungen

SymptomUrsacheLösung
Es kommen keine Werte anbroker sendet keine RecordsMit kafka-console-consumer.sh --bootstrap-server <host>:9092 --topic <t> --from-beginning direkt empfangen und prüfen
[KAFKA] connect 실패: TimeoutExceptionbootstrap-Server nicht erreichbartelnet <host> 9092 prüfen. Kontrollieren, ob advertised.listeners auf die externe IP zeigt (broker-seitig)
[KAFKA] poll error: Authentication failedSASL-Authentifizierung fehlgeschlagensasl-mechanism / sasl-jaas-config erneut prüfen. Sonderzeichen im Passwort escapen
[KAFKA] poll error: SslHandshakeTLS-Truststore fehlt / abgelaufenbroker-Zertifikat von vertrauenswürdiger CA ausstellen lassen oder zum System-Truststore hinzufügen
Neue Group empfängt keine Datenauto-offset=latest gesetzt, aber keine neuen Records vorhandenEinmalig auf earliest umstellen oder warten, bis ein Publish erfolgt

Weiterführende technische Dokumentation