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
- 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:
pip install connic-composer-sdk
connic init my-project --templates=kafka-fraud-detector
cd my-projectDadurch entsteht eine vollständige, deploybare Struktur:
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.txtDie 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:
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 == TrueBei 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:
connic login
connic deployWä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:
| Einstellung | Wert | Hinweise |
|---|---|---|
| Mode | Inbound (Consumer) | Konsumiert Nachrichten und startet Runs |
| Bootstrap Servers | kafka:9092 | Kommagetrennte Broker-Adressen, von Connic aus erreichbar |
| Topic | transactions | Das zu konsumierende Topic |
| Consumer Group ID | fraud-detector | Optional; wird bei leerem Feld automatisch generiert |
| Auto Offset Reset | latest | Mit earliest vom Anfang erneut abspielen |
| Security Protocol | PLAINTEXT | Oder SSL, SASL_PLAINTEXT, SASL_SSL |
| SASL Mechanism | PLAIN, SCRAM-SHA-256 oder SCRAM-SHA-512 | Zusä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:
{
"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:
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.
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 installierenDie Verbindung funktioniert auch für Streams zur Inhaltsmoderation, Personalisierung und Anomalieerkennung. Der Marketplace beschreibt den Funktionsumfang der Kafka-Verbindung.