Apache Kafka
Ein Agent abonniert ein Topic und verarbeitet jede Nachricht als Run. Sein Output lässt sich für nachgelagerte Consumer in einem anderen Topic veröffentlichen. Consumer Groups, Offsets und Reconnects werden verwaltet.
Überblick
Die Kafka-Verbindung ist ein verwalteter Consumer und Producer für Agenten. Eine eingehende Verbindung mit den Brokern und einem Topic macht jede Nachricht darin zu einem Agent Run: Bei JSON-Objekten bleiben die Felder auf oberster Ebene erhalten, alles andere landet unter dem Key message. Jeder Run enthält außerdem Topic, Partition, Offset, Timestamp und Message Key als _kafka Metadaten.
Der sonst nötige Consumer Service entfällt. Connic verwaltet die Consumer Group, verfolgt Offsets, verbindet sich mit Exponential Backoff neu und übernimmt Konfigurationsänderungen. Eine ausgehende Verbindung schreibt abgeschlossene Ergebnisse in ein Ziel-Topic. Eine vollständige Consume-Process-Produce-Pipeline braucht dann keinen weiteren Code außer dem Agenten selbst.
So funktioniert es
Inbound Consumer erstellen
Füge im Verbindungsdiagramm des Agenten eine Apache-Kafka-Verbindung im Modus Inbound (Consumer) hinzu. Trage Bootstrap Servers und das zu lesende Topic sowie bei Bedarf Consumer Group ID und Auto Offset Reset ein. latest liest nur neue Nachrichten, earliest beginnt am Anfang.
Sicherheit konfigurieren
Wähle PLAINTEXT, SSL, SASL_PLAINTEXT oder SASL_SSL. Verwende für Managed Kafka wie Confluent, MSK oder Aiven SASL_SSL mit SCRAM-SHA-256. Füge für TLS und Mutual TLS Zertifikat und Schlüssel im PEM-Format direkt in das Formular ein.
Agenten verknüpfen und optional Ergebnisse veröffentlichen
Jede Nachricht im Topic wird an alle verknüpften Agenten gesendet. Eine ausgehende Verbindung mit Bootstrap Servers und Ziel-Topic veröffentlicht abgeschlossene Run-Outputs automatisch im nachgelagerten System.
Muster für den Produktivbetrieb
Muster für den Produktivbetrieb, bei denen Connic Queues, Worker und Scheduler betreibt.
Nachricht veröffentlicht
topic: user.content.createdDer Agent bewertet jeden Post auf Toxizität und Richtlinienkonformität und veröffentlicht, prüft oder verbirgt ihn, bevor andere Nutzer den Kommentar sehen.
Nachricht veröffentlicht
topic: user.content.createdDer Agent bewertet jeden Post auf Toxizität und Richtlinienkonformität und veröffentlicht, prüft oder verbirgt ihn, bevor andere Nutzer den Kommentar sehen.
Nutzer klickt auf einen Inhalt
topic: user.clickstreamDer Agent aktualisiert das Interessenprofil, prognostiziert das nächste gewünschte Produkt und schreibt eine neue Empfehlung in ein Ergebnis-Topic, das das Frontend liest.
Nutzer klickt auf einen Inhalt
topic: user.clickstreamDer Agent aktualisiert das Interessenprofil, prognostiziert das nächste gewünschte Produkt und schreibt eine neue Empfehlung in ein Ergebnis-Topic, das das Frontend liest.
Transaktion gestartet
topic: payments.transactionsDer Agent gleicht die Transaktion mit Verhaltensmustern, Device Fingerprints und der Historie ab und veröffentlicht vor Abschluss der Autorisierung eine Entscheidung zur Freigabe oder Ablehnung.
Transaktion gestartet
topic: payments.transactionsDer Agent gleicht die Transaktion mit Verhaltensmustern, Device Fingerprints und der Historie ab und veröffentlicht vor Abschluss der Autorisierung eine Entscheidung zur Freigabe oder Ablehnung.
Sensordaten treffen ein
topic: iot.sensor.readingsDer Agent markiert Werte außerhalb der Spezifikation, prognostiziert das ausfallende Bauteil und eröffnet vor dem Stillstand der Linie einen Arbeitsauftrag mit Diagnose.
Sensordaten treffen ein
topic: iot.sensor.readingsDer Agent markiert Werte außerhalb der Spezifikation, prognostiziert das ausfallende Bauteil und eröffnet vor dem Stillstand der Linie einen Arbeitsauftrag mit Diagnose.
Message-Payload und Consumer-Semantik
Bei JSON-Objekten werden die Felder auf oberster Ebene als Payload übergeben, ergänzt um einen _kafka Block mit Topic, Partition, Offset, Timestamp und Key. Nicht-JSON-Werte landen unter message, und Compaction Tombstones lösen weiterhin Runs mit message: null aus. So kann der Agent über _kafka.key auf Löschungen reagieren. Consumer Groups verhalten sich wie gewohnt: Verbindungen mit unterschiedlichen Group IDs sehen jeweils jede Nachricht, Verbindungen mit derselben Group ID teilen die Last.
Jede konsumierte Nachricht wird zu einem normalen Agent Run mit vollständigen Traces, Token- und Kosten-Tracking sowie denselben Guardrails und Approval-Regeln wie jeder andere Trigger. Auch ein Topic mit hohem Durchsatz wird nicht zur Black Box.
Ergebnisse zurückschreiben
Ausgehende Verbindungen veröffentlichen Ergebnisse, sobald verknüpfte Agenten abgeschlossen sind. Nur Runs mit dem Status completed werden veröffentlicht; fehlgeschlagene und abgebrochene Runs werden übersprungen. Wenn eine Inbound-Nachricht einen Key hatte, wird das Ergebnis mit demselben Key erzeugt. Dadurch bleibt die Reihenfolge innerhalb der Partition über die gesamte Pipeline erhalten; ohne Inbound-Key wird die Run ID verwendet. Nachrichten werden mit vollständiger Replikationsbestätigung und automatischen Wiederholungen veröffentlicht. Kafka-Dokumentation: Payloads und vollständige Pipeline.
Informationen
- Publisher
- Von Connic
- Verbindungen
- Verbindungen
- Modi
- Inbound, Outbound
- Dokumentation
- Apache Kafka Docs
Häufig gestellte Fragen
_kafka. Ein vollständiges Beispiel zeigt die Kafka-Pipeline-Anleitung.Beschreibe die Event-Quelle, die Struktur der Payload, das Ziel für die Ergebnisse und die Anforderungen an private Netzwerke oder Freigaben. Wir helfen, Apache Kafka mit dem passenden Verbindungsmodus, Deployment und Monitoring in Connic einzurichten.
Sprich mit Sales