Zum Hauptinhalt springen
Connic

KI-Agenten über Kafka-Topics auslösen

Eine eingehende Connic-Kafka-Verbindung auf einem Topic startet mit jeder Nachricht einen Agenten-Run. Kafka-Verbindung konfigurieren, Agenten verknüpfen, bereitstellen und Runs beobachten.

12. Juli 2026(zuletzt aktualisiert: 19. Juli 2026)8 Min. LesezeitAutor: Connic Engineering

Transaktionen, Bestellungen und Telemetriedaten fließen bereits durch Kafka. Was fehlt, ist ein KI-Agent, der darauf reagiert: die Transaktion bewertet, den Alert priorisiert oder die Empfehlung aktualisiert. Connic benötigt dafür keinen separaten Consumer-Service. Eine eingehende Kafka-Verbindung wird mit dem Topic verbunden, und jede Nachricht startet einen Agenten-Run mit den Daten der Nachricht. Das Tutorial führt vom Grundgerüst des Projekts bis zum ersten beobachteten Run.

Warum Kafka ein starker Agenten-Trigger ist

Kafka ist ein dauerhaftes, geordnetes Log. Das Triggern von Agenten daraus bietet At-Least-Once Delivery, eine feste Reihenfolge innerhalb einer Partition und die Möglichkeit, den Verlauf durch das Zurücksetzen eines Consumers auf einen früheren Offset erneut abzuspielen. Ein einfacher HTTP Webhook kann diese Garantien nicht bieten. Die Verbindung tritt dem Cluster als Consumer bei. Zwischen Kafka und den Agenten muss daher weder ein Polling-Loop noch ein Middleware- oder Bot-Prozess dauerhaft laufen. Bei der Wahl des passenden Triggers hilft der Vergleich von Webhook-, Kafka-, Postgres- und SQS-Triggern. Hier geht es um die praktische Umsetzung mit Kafka.

Schritt 1: Ein Projekt aus dem Template aufsetzen

Voraussetzungen
  • Installiertes Python 3.10+
  • Ein erreichbarer Kafka-Cluster mit einem zu konsumierenden Topic (dieses Tutorial verwendet transactions)
  • Ein mit einem Git-Repository verbundenes Connic-Projekt
  • Das Template legt ein Anthropic-KI-Modell fest. Hinterlege daher entweder einen Anthropic-Key in den Projekteinstellungen oder stelle die Agent-YAML auf ein exaktes connic/* KI-Modell um und verwende Projektguthaben

Der schnellste Einstieg ist das Kafka Fraud Detector Template aus dem Connic Marketplace. Das funktionsfähige Projekt verarbeitet Transaktionen aus einem Kafka-Topic, bewertet deren Betrugsrisiko und enthält einen Eskalationsagenten für Hochrisikofälle. Installiere das SDK und erstelle daraus das Grundgerüst:

terminal
pip install connic-composer-sdk
connic init my-project --templates=kafka-fraud-detector
cd my-project

Dadurch entsteht eine vollständige, deploybare Struktur:

Dateistruktur
my-project/
  agents/
    kafka-fraud-detector/
      fraud-scorer.yaml
      fraud-escalator.yaml
  tools/
    fraud_tools.py
  middleware/
    fraud-scorer.py
  hooks/
    fraud-scorer.py
  schemas/
    fraud-assessment.json
  tests/
  requirements.txt

Die Struktur enthält zwei Agenten, reine Python-Tools, Middleware zur Kennzeichnung von Admin-Requests, einen Hook zur Durchsetzung des Admin-Overrides, ein Output-Schema und eine Test-Suite, die mit connic test ausgeführt wird. Das Kafka Fraud Detector Template enthält ein Architekturdiagramm. Für andere Anwendungsfälle gibt es weitere Agentenvorlagen im Marketplace. Alles Folgende funktioniert für einen von Grund auf selbst geschriebenen Agenten genauso.

Schritt 2: Die Agenten-YAML ansehen

Jede YAML-Datei unter agents/ definiert einen Agenten. Hier folgt der Scorer aus dem Template mit gekürztem System-Prompt:

agents/kafka-fraud-detector/fraud-scorer.yaml
version: "1.0"

name: fraud-scorer
type: llm
model: connic/gpt-5.6-terra
description: "Bewertet Transaktionen aus einem Kafka Stream in Echtzeit auf Betrugsrisiken"

system_prompt: |
  Du bist ein Spezialist für Betrugserkennung und analysierst Finanztransaktionen
  aus einem Kafka Stream in Echtzeit.

  Verwende _kafka.key zur Kundenzuordnung und _kafka.timestamp, um
  Verarbeitungslatenzen zu erkennen.

temperature: 0.1
output_schema: fraud-assessment

concurrency:
  key: "data.customer_id"
  on_conflict: queue

tools:
  - fraud_tools.calculate_velocity
  - fraud_tools.check_geo_anomaly
  - fraud_tools.create_alert
  - fraud_tools.search_fraud_patterns
  - fraud_tools.store_fraud_pattern
  - fraud_tools.admin_override: context.is_admin == True

Bei Kafka-Workloads weist der Prompt das KI-Modell an, _kafka.key und _kafka.timestamp zu verwenden. Diese Metadaten hängt die Verbindung an jede Nachricht. Der concurrency-Block stellt Runs mit derselben customer_id in eine Queue. Dadurch konkurrieren zwei Transaktionen desselben Kunden nie miteinander, während der Rest des Streams parallel verarbeitet wird. Der letzte Tool-Eintrag ist außerdem bedingt: admin_override steht dem KI-Modell nur zur Verfügung, wenn die Middleware im Run-Kontext is_admin gesetzt hat. Die Dokumentation zur Agentenkonfiguration beschreibt die Funktionsweise und alle verfügbaren Optionen.

Schritt 3: Agenten bereitstellen

Pushe das Projekt auf den Deployment-Branch der Umgebung. Connic erstellt einen Build und stellt das Projekt automatisch bereit. Alternativ lässt es sich direkt über die CLI bereitstellen:

terminal
connic login
connic deploy

Während der Iteration stellt connic dev einen Dev-Server mit Hot Reload bereit. Prompt und Tools lassen sich damit anpassen, bevor produktiver Traffic angebunden wird.

Schritt 4: Die eingehende Kafka-Verbindung erstellen

Verbindungen werden im Dashboard erstellt und dort mit Agenten verknüpft. Öffne die Detailseite des Agenten fraud-scorer, klicke im Verbindungsablauf auf Add inbound connector, wähle Create New Connector und anschließend Apache Kafka. Für das Tutorial gilt:

EinstellungWertHinweise
ModeInbound (Consumer)Konsumiert Nachrichten und startet Runs
Bootstrap Serverskafka:9092Kommagetrennte Broker-Adressen, von Connic aus erreichbar
TopictransactionsDas zu konsumierende Topic
Consumer Group IDfraud-detectorOptional; wird bei leerem Feld automatisch generiert
Auto Offset ResetlatestMit earliest vom Anfang erneut abspielen
Security ProtocolPLAINTEXTOder SSL, SASL_PLAINTEXT, SASL_SSL
SASL MechanismPLAIN, SCRAM-SHA-256 oder SCRAM-SHA-512Zusätzlich Username und Password, nur bei SASL

Die Verbindung empfängt Nachrichten über unsere Infrastruktur. Die Broker müssen daher von Connic aus direkt oder über Bridge erreichbar sein: Ein verwalteter Endpoint funktioniert direkt, ein privater Cluster wird über Bridge verbunden, wie am Ende dieses Artikels beschrieben. Für Managed Kafka wie Confluent Cloud, Amazon MSK oder Aiven werden SASL_SSL mit SCRAM-SHA-256 und die Service-Zugangsdaten verwendet. Mutual TLS wird durch direktes Einfügen der PEM-Inhalte von CA Certificate, Client Certificate und Client Key in das Formular unterstützt. Die Verbindung akzeptiert keine lokalen Dateipfade für Zertifikate. Nach der Erstellung wird sie automatisch mit dem Agenten verknüpft und beginnt mit dem Konsumieren.

Consumer Groups verhalten sich wie von Kafka gewohnt: Verbindungen mit unterschiedlichen Group IDs erhalten jeweils jede Nachricht; Verbindungen mit derselben Group ID teilen die Last. Die Verbindung verfolgt Offsets, stellt die Verbindung mit Exponential Backoff wieder her und übernimmt Konfigurationsänderungen, ohne den Agenten erneut bereitzustellen.

Was der Agent erhält

JSON-Nachrichten kommen mit unveränderten Top-Level-Feldern an. Hinzu kommt ein _kafka-Block mit Topic, Partition, Offset, Timestamp und Key:

message-payload.json
{
  "transaction_id": "tx-9917",
  "customer_id": "cust-1042",
  "amount": 4899.0,
  "currency": "EUR",
  "merchant": "TechWorld Berlin",
  "country": "DE",
  "_kafka": {
    "topic": "transactions",
    "partition": 0,
    "offset": 1542,
    "timestamp": 1705312800000,
    "key": "cust-1042"
  }
}

Nicht-JSON-Nachrichten werden unter einem message-Key verpackt. Compaction-Tombstones lösen weiterhin Runs mit message: null aus, sodass ein Agent über _kafka.key auf Löschungen reagieren kann. Das vollständige Payload-Verhalten beschreibt die Kafka-Verbindungsreferenz.

Schritt 5: Ergebnisse in ein Topic zurückschreiben

Nachgelagerte Systeme benötigen die Entscheidung des Agenten meist in einem eigenen Stream. Wiederhole die Einrichtung der Verbindung für denselben Agenten fraud-scorer. Wähle diesmal den Modus Outbound (Producer) und fraud-alerts als Topic. Nach jedem abgeschlossenen Run veröffentlicht der Agent seine strukturierte Bewertung zusammen mit Run-ID, Agentenname, Status und Output. Fehlgeschlagene und abgebrochene Runs werden übersprungen. Wurde der Run durch eine eingehende Kafka-Nachricht mit Key ausgelöst, verwendet die ausgehende Nachricht denselben Key. So bleibt die Reihenfolge pro Kunde über den gesamten Roundtrip durch den Agenten erhalten.

Schritt 6: Ein Test-Event senden und den Run beobachten

Sende eine Testnachricht mit Key an das Topic:

terminal
kafka-console-producer \
  --bootstrap-server kafka:9092 \
  --topic transactions \
  --property parse.key=true \
  --property key.separator=:
>cust-1042:{"transaction_id":"tx-9917","customer_id":"cust-1042","amount":4899.0,"currency":"EUR","merchant":"TechWorld Berlin","country":"DE"}

Kurz darauf erscheint im Runs-Bereich des Projekts ein neuer Run mit dem vollständigen Trace: dem eingehenden Payload, jedem Tool-Call des Scorers und seiner strukturierten Risikobewertung. Die Trace-Dokumentation erklärt die Auswertung. Danach ist der Feedback-Loop kurz: Passe Prompt oder Tools an, pushe, und die nächste Nachricht im Topic führt die neue Version aus.

Private Kafka-Cluster

Broker in einem privaten Netzwerk müssen nicht öffentlich zugänglich sein. Mit der Option Connect via Bridge und einer im privaten Netzwerk laufenden Connic Bridge nutzt die Verbindung einen sicheren ausgehenden Tunnel.

Starte mit einer funktionierenden Kafka-Pipeline

Das Kafka Fraud Detector Template enthält zwei Agenten, Tools, Middleware, Schemas und Tests. Erstelle das Grundgerüst mit einem Befehl und ersetze die Beispiellogik durch deinen Anwendungsfall.

Kafka Fraud Detector Template installieren

Die Verbindung funktioniert auch für Streams zur Inhaltsmoderation, Personalisierung und Anomalieerkennung. Der Marketplace beschreibt den Funktionsumfang der Kafka-Verbindung.

Häufig gestellte Fragen

In Connic wird eine eingehende Kafka-Verbindung mit den Bootstrap-Servern und optionalen SASL- oder SSL-Einstellungen mit dem Topic verbunden und einem bereitgestellten Agenten zugeordnet. Die Verbindung tritt dem Topic als Consumer bei und startet für jede Nachricht einen Agenten-Run mit den Daten der Nachricht und Kafka-Metadaten. Consumer-Code ist nicht erforderlich.

Ja. Mit Auto Offset Reset auf earliest verarbeitet eine neue Consumer Group das Topic von Anfang an. So lässt sich der Verlauf nach der Verbesserung eines Agenten erneut verarbeiten. Die Zustellung erfolgt At-Least-Once; der Agent sollte daher idempotent sein, etwa durch Deduplizierung anhand des Kafka-Message-Keys.

Ja. Eine ausgehende Kafka-Verbindung wird hinzugefügt und mit dem Agenten verknüpft. Abgeschlossene Runs werden mit Run-ID, Agentenname, Status und Output im konfigurierten Topic veröffentlicht; fehlgeschlagene und abgebrochene Runs werden übersprungen. Ergebnisse, die durch eine eingehende Nachricht mit Key ausgelöst wurden, verwenden denselben Key und bewahren dadurch die Reihenfolge in der Partition.

Ja. SASL_SSL dient als Security Protocol mit dem Mechanism SCRAM-SHA-256 und den Service Credentials. Für Mutual TLS werden die PEM-Inhalte von CA Certificate, Client Certificate und Client Key direkt in das Verbindungsformular eingefügt; lokale Dateipfade werden nicht unterstützt.

Mehr aus dem Blog

Tutorial

Python-KI-Agenten ohne Kubernetes bereitstellen

Ein Python-KI-Agent lässt sich mit YAML, reinem Python, Tests als Deployment-Gate, Git und einer Managed Runtime in der EU ohne Kubernetes bereitstellen. Der Artikel enthält funktionsfähigen Code.

12. August 202612 Min. Lesezeit
Tutorial

So integrieren kleine Engineering-Teams einen KI-Agenten in SaaS

Ein praktikabler Weg zum ersten KI-Agenten im Produktivbetrieb: Aufgabe eingrenzen, per Konfiguration definieren, vorhandene Systeme anbinden und die Runtime verwalten lassen.

12. Juni 20269 Min. Lesezeit
Tutorial

KI-Agenten automatisch mit LLM Judges bewerten

Ein LLM Judge bewertet ausgewählte oder alle passenden Agenten-Runs nach festgelegten Kriterien. Score-Trends und Alerts machen Regressionen sichtbar.

29. März 202610 Min. Lesezeit
Tutorial

LangChain-KI-Agenten in den Produktivbetrieb migrieren

Ein funktionierender LangChain-Prototyp muss echten Traffic bewältigen. Bestehender Agenten-Code lässt sich ohne vollständige Neuentwicklung auf eine Plattform für den Produktivbetrieb migrieren.

23. März 202611 Min. Lesezeit
Tutorial

Datenbank, Retrieval oder Sessions: Speicher für Agenten im Vergleich

Connics Datenbank, Retrieval und persistente Sessions im Vergleich: passende Einsatzbereiche, gespeicherte Gesprächsverläufe, TTL und zugehörige Runs.

4. März 202612 Min. Lesezeit
Tutorial

KI-Agenten: Vom Prototyp zum Produktivbetrieb

Eine Demo funktioniert hervorragend, bis 1.000 Nutzer gleichzeitig darauf zugreifen. Dieser Leitfaden beschreibt die oft zu spät berücksichtigten Anforderungen des Produktivbetriebs.

10. Januar 202610 Min. Lesezeit
Tutorial

Versteckte Kosten beim Self-Hosting von KI-Agenten

Ein Kubernetes-Deployment wirkt zunächst einfach. Der Artikel vergleicht die tatsächlichen Kosten selbst gehosteter KI-Agenten mit einer verwalteten Plattform.

18. Dezember 20257 Min. Lesezeit
Tutorial

KI-Agenten ohne ML-Team in SaaS integrieren

Kunden erwarten KI-Features, auch wenn das Produktteam keine ML-Engineers beschäftigt. Vorhandene Fähigkeiten reichen aus, um KI-Agenten zu veröffentlichen.

5. Dezember 20258 Min. Lesezeit
Tutorial

RAG-Tutorial für KI-Agenten: Retrieval mit Quellenangaben

Ein RAG-Agent für den Produktivbetrieb braucht abgegrenzte Retrieval-Namespaces, Read-only-Berechtigungen, Quellenangaben, eigene Tool-Wrapper und Regressionstests.

15. November 20259 Min. Lesezeit