Kafka
Apache Kafka verbindet Agenten mit Event Streams in Echtzeit. Eingehende Nachrichten starten Agenten; ausgehende Verbindungen veröffentlichen deren Ergebnisse in Topics.
Auf dieser Seite
Einrichtung
Kafka-Cluster vorbereiten
Stelle sicher, dass der Cluster läuft und das Quell-Topic existiert.
Verbindung erstellen
Öffne den Agenten, klicke im Verbindungsflussdiagramm auf Add inbound connector, dann auf Create New Connector und wähle Apache Kafka.
Konfigurieren und erstellen
Wähle den Mode Inbound (Consumer), gib Bootstrap Servers und Topic sowie optional Group ID, Auto Offset Reset und Sicherheitseinstellungen (SASL/SSL) ein. Klicke auf Create; Connic beginnt anschließend, Nachrichten aus dem Topic zu konsumieren.
So funktioniert Inbound
Inbound-Kafka-Verbindungen arbeiten als Consumer, die ein Topic abonnieren. Trifft eine Nachricht ein, wird sie geparst und an alle verknüpften Agenten weitergeleitet. Das eignet sich für ereignisgesteuerte Verarbeitung, Echtzeit-Pipelines, Streaming Analytics und die Kommunikation zwischen Microservices.
- Bootstrap Servers: Adressen der Kafka-Broker (z. B.
kafka:9092) - Topic: Topic, aus dem Nachrichten konsumiert werden
- Group ID (optional): Kennung der Consumer Group
- Auto Offset Reset (optional): Wenn kein bestätigter Offset existiert, liest Latest Nachrichten, die nach dem Start des Consumers erzeugt wurden, während Earliest beim ersten verfügbaren Offset beginnt
Sicherheitseinstellungen
- Security Protocol: PLAINTEXT, SSL, SASL_PLAINTEXT oder SASL_SSL (empfohlen)
- SASL Mechanism: PLAIN, SCRAM-SHA-256 oder SCRAM-SHA-512
- Username & Password: Für SASL erforderlich
- SSL CA Certificate PEM (optional): Inhalt des CA-Zertifikats für die TLS-Verifizierung
- SSL Client Certificate PEM & Client Key PEM (optional): Inhalt des Client-Zertifikats und privaten Schlüssels für Mutual TLS; gib beide Felder gemeinsam an
- SSL Key Password (optional): Passwort für einen verschlüsselten privaten Client-Schlüssel
- Verify SSL Hostname (optional): Prüft die Hostnamen der Broker gegen das TLS-Zertifikat
Füge die PEM-Inhalte des Zertifikats und Schlüssels direkt in das Formular ein. Die Kafka-Verbindungskonfiguration akzeptiert keine lokalen Dateipfade für Zertifikate oder Schlüssel.
Tipp: Verwende für Managed Kafka (Confluent, MSK, Aiven) SASL_SSL mit SCRAM-SHA-256.
Nachrichten-Payload
{
"order_id": "12345",
"customer": "john@example.com",
"items": ["widget-a", "widget-b"],
"_kafka": {
"topic": "orders",
"partition": 0,
"offset": 1542,
"timestamp": 1705312800000,
"key": "order-12345"
}
}Die _kafka-Metadaten enthalten topic, partition, offset, timestamp und key.
Nachrichten mit JSON-Objekten werden wie oben gezeigt mit ihren Feldern auf oberster Ebene weitergeleitet. Alle anderen Werte werden unter einem message-Schlüssel verpackt: Nicht-JSON-Werte kommen als {"message": "<raw text>", "_kafka": ...} an, und Nachrichten mit null-Wert (Compaction Tombstones) starten weiterhin Runs mit message: null. Verwende _kafka.key, um die gelöschte Entität zu identifizieren:
{
"message": null,
"_kafka": {
"topic": "orders",
"partition": 0,
"offset": 1543,
"timestamp": 1705312800000,
"key": "order-12345"
}
}End-to-End-Beispiel
Erstelle eine Inbound-Verbindung für das Quell-Topic, verknüpfe sie mit einem Agenten und verknüpfe außerdem eine Outbound-Verbindung, die die Agentenausgabe in einem Ziel-Topic veröffentlicht.
version: "1.0"
name: order-processor
type: llm
model: connic/gpt-5.6-luna
description: "Validate orders and compute routing"
system_prompt: |
You receive an order event in JSON (from Kafka).
1) Call orders.validate_order
2) Call orders.score_risk
3) Return JSON with order_id, status, risk_score, route
tools:
- orders.validate_order
- orders.score_risk
output_schema: order-result.jsonfrom typing import Dict, Any
def validate_order(order_id: str, items: list[str]) -> Dict[str, Any]:
"""Basic order validation."""
if not order_id or not items:
return {"ok": False, "reason": "missing_fields"}
return {"ok": True}
async def score_risk(customer: str, total: float) -> Dict[str, Any]:
"""Return a simple risk score and routing hint."""
score = 0.02 if total < 100 else 0.12
route = "standard" if score < 0.1 else "manual_review"
return {"risk_score": score, "route": route}Die JSON-Ausgabe des Agenten erscheint im Feld output des Outbound-Run-Wrappers:
{
"order_id": "12345",
"status": "approved",
"risk_score": 0.02,
"route": "standard"
}Consumer Groups
Jede Verbindung verwendet eine Consumer Group, um verarbeitete Nachrichten zu verfolgen. Unterschiedliche Group IDs = dieselben Nachrichten für alle. Gleiche Group ID = Lastverteilung.
Ausfallsichere Verbindung
- Stellt die Verbindung mit exponentiellem Backoff wieder her und bestätigt Offsets nach erfolgreicher Verarbeitung
Einrichtung
Kafka-Cluster vorbereiten
Stelle sicher, dass der Cluster läuft und das Ziel-Topic existiert, oder aktiviere die automatische Erstellung auf dem Broker.
Verbindung erstellen
Öffne den Agenten, klicke auf Add outbound connector, dann auf Create New Connector und wähle Apache Kafka.
Konfigurieren und erstellen
Wähle den Mode Outbound (Producer) und gib Bootstrap Servers sowie Topic ein. Klicke auf Create; Ergebnisse verknüpfter Agenten werden anschließend im Topic veröffentlicht.
So funktioniert Outbound
Automatische Outbound-Verbindungen veröffentlichen die Wrapper abgeschlossener Runs im konfigurierten Topic und können auf ausgewählte Inputs begrenzt werden. Agent-Tool- und Middleware-Outbound-Verbindungen veröffentlichen nur bei einem Aufruf.
Nur Runs mit dem Status completed werden veröffentlicht. Fehlgeschlagene und abgebrochene Runs werden übersprungen. Runs, die durch StopProcessing vorzeitig beendet wurden, schließen regulär ab und werden daher ebenfalls veröffentlicht; die StopProcessing-Antwort dient dabei als Ausgabe. Löse StopProcessing mit publish_outbound=False aus, um die automatische Outbound-Verbindung für diesen Run zu überspringen. Ein bereits erfolgter Aufruf einer Agent-Tool- oder Middleware-Outbound-Verbindung wird dadurch nicht rückgängig gemacht.
- Bootstrap Servers: Adressen der Kafka-Broker
- Topic: Topic, in dem Ergebnisse veröffentlicht werden
Sicherheit
Wie bei inbound: Für Managed Kafka Services wird SASL_SSL mit SCRAM-SHA-256 empfohlen. Füge für Mutual TLS die PEM-Inhalte in die Felder SSL CA Certificate PEM, SSL Client Certificate PEM und SSL Client Key PEM ein; lokale Dateipfade zu Zertifikaten oder Schlüsseln werden nicht unterstützt.
Automatische Payload
{
"run_id": "550e8400-e29b-41d4-a716-446655440000",
"agent_name": "order-processor",
"status": "completed",
"output": "Order processed successfully. Total: $234.56",
"error": null,
"started_at": "2024-01-15T10:30:00Z",
"ended_at": "2024-01-15T10:30:05Z",
"token_usage": {
"input_tokens": 150,
"output_tokens": 50,
"thinking_tokens": 0,
"cached_input_tokens": 0,
"total_tokens": 200
}
}Enthält run_id, agent_name, status, output, error, Zeitstempel und token_usage.
Agent-Tool- und Middleware-Outbound-Verbindungen
Eine Agent-Tool-Outbound-Verbindung stellt einen bearbeitbaren Tool-Namen bereit, standardmäßig send_to_<connector_name>. Rufe eine Middleware-Outbound-Verbindung über send_connector mit ihrem konfigurierten Namen auf. Beide verwenden dieses von der Verbindung definierte Payload-Schema:
{
"payload": {
"order_id": "12345",
"status": "approved"
},
"key": "order-12345"
}payload wird zum Wert der Kafka-Nachricht. key ist optional. Connic verwendet das konfigurierte Topic, die Verbindung, Serialisierung, Sicherheit, Wiederholungsstrategie und das Bridge-Routing, ohne dem Modell Zugangsdaten offenzulegen.
Korrelation über den Message Key
Bei einer automatischen Outbound-Verbindung verwendet ein durch eine Inbound-Kafka-Verbindung gestarteter Run erneut den Message Key der Quelle; andere Runs verwenden die Run-ID. Agent-Tool- und Middleware-Outbound-Verbindungen können key ausdrücklich angeben.
# Eingehende Nachricht mit dem Key "order-123"
Nachricht empfangen → Agent ausgelöst → Run abgeschlossen
# Ausgehende Nachricht verwendet denselben Key "order-123"
Ergebnis mit Key veröffentlicht: "order-123" (Quelle: ursprüngliche Nachricht)
# Ohne eingehenden Key wird run_id verwendet
Ergebnis mit Key veröffentlicht: "550e8400-e29b..." (Quelle: run_id)Zustellungsgarantien
- Bestätigung nach vollständiger Replikation und automatische Wiederholungsversuche
Fortgeschrittene Muster
Erstelle für Multi-Topic-Streams pro Topic eine Inbound-Verbindung und verknüpfe alle mit demselben Korrelations-Agenten. Verwende _kafka.key (oder eine gemeinsame order_id), um Ereignisse zu korrelieren und Teilzustände über Tools in einer externen Datenbank oder Redis zu speichern. Wiederholungsversuche werden am Agenten konfiguriert:
version: "1.0"
name: event-correlator
type: llm
model: connic/gpt-5.6-luna
description: "Correlate order + shipment events across topics"
system_prompt: |
Use the correlation tools to store incoming events by _kafka.key.
Only respond when both "orders" and "shipments" are present.
tools:
- correlation.upsert_event
- correlation.build_snapshot
retry_options:
attempts: 5
initial_delay: 10
max_delay: 30