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.
| Situation | Welcher Modus ist zu verwenden |
|---|---|
| Ein externes IT-System sendet Werte per REST POST | HTTP-Push |
| Ein externer Server pusht per WebSocket | WebSocket Client |
| Ein externer IIoT-MQTT broker pusht auf ein Topic | MQTT Client |
| Kafka topic wird als Consumer empfangen | Apache Kafka (diese Seite) |
| Werte werksinterner PLCs werden direkt gelesen | Modbus / OPC-UA usw. |
Eingabefelder des Registrierungsformulars
| Eingabefeld | Was wird eingetragen | Beispiel |
|---|---|---|
| IP-Adresse | Kafka-broker-Host (bootstrap) | kafka.example.com, 10.0.0.50 |
| Port | broker-Port (unverschlüsselt 9092, TLS 9093) | 9092, 9093 |
| GROUP ID | consumer group id (optional) | plant-floor-A, edge-line1 |
| SECURITY | Sicherheitsprotokoll | PLAINTEXT / SSL / SASL_PLAINTEXT / SASL_SSL |
| SASL MECH | SASL-Mechanismus (optional) | PLAIN / SCRAM-SHA-256 / SCRAM-SHA-512 |
| SASL JAAS | JAAS-Konfiguration für die SASL-Authentifizierung (optional) | org.apache.kafka.common.security.plain.PlainLoginModule required username="u" password="p"; |
| AUTO OFFSET | Startpunkt für neue Groups | latest (Standard) / earliest |
| Erfassungsintervall | Intervall, 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.
| Modus | Format | Verhalten |
|---|---|---|
| SCALAR | sensors/temp oder sensors/temp.value | Gesamte Nachricht als String. Ist die Nachricht ein JSON-Objekt, erfolgt ein Raw-Fallback. |
| KEY | sensors/temp:temperature | Wert eines Top-Level-JSON-Keys (z. B. {"temperature":25.3,"humidity":60} → 25.3) |
| PATH | sensors/temp:$.data.tags.T1 | Dynamische Auswertung per JSON Pointer (verschachtelte Keys unterstützt) |
| RAW | sensors/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
| Anwendungsfall | Vorgehen |
|---|---|
| Werksinternes Kafka-Cluster | host = bootstrap-broker-IP, 9092 (unverschlüsselt) |
| Confluent Cloud | host = 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 + earliest | AUTO OFFSET=earliest wählen, um alle vorhandenen Daten zu empfangen |
| Getrennte Verwaltung der Group-Offsets | Mit 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
| Symptom | Ursache | Lösung |
|---|---|---|
| Es kommen keine Werte an | broker sendet keine Records | Mit kafka-console-consumer.sh --bootstrap-server <host>:9092 --topic <t> --from-beginning direkt empfangen und prüfen |
[KAFKA] connect 실패: TimeoutException | bootstrap-Server nicht erreichbar | telnet <host> 9092 prüfen. Kontrollieren, ob advertised.listeners auf die externe IP zeigt (broker-seitig) |
[KAFKA] poll error: Authentication failed | SASL-Authentifizierung fehlgeschlagen | sasl-mechanism / sasl-jaas-config erneut prüfen. Sonderzeichen im Passwort escapen |
[KAFKA] poll error: SslHandshake | TLS-Truststore fehlt / abgelaufen | broker-Zertifikat von vertrauenswürdiger CA ausstellen lassen oder zum System-Truststore hinzufügen |
| Neue Group empfängt keine Daten | auto-offset=latest gesetzt, aber keine neuen Records vorhanden | Einmalig auf earliest umstellen oder warten, bis ein Publish erfolgt |