Playbook ┬╖ 6 minute read
How to Build a Kafka AI Pipeline
A Kafka AI pipeline consumes events, enriches them with model inference, and produces results to downstream topics, designed around the fact that inference is far slower than the stream. That means batching, bounded concurrency, backpressure rather than unbounded lag, idempotent processing keyed by event identity, versioned schemas, and routing so only events that need a model reach one.
Kafka moves events at a rate no model can match. A topic delivering thousands of events a second to a consumer that calls a model taking hundreds of milliseconds per event falls behind immediately and never recovers. That mismatch is the central fact of streaming AI, and every design decision in this guide follows from it. It draws on FISTA Solutions' AI agents delivery on event-driven platforms and complements kafka vs rabbitmq for ai pipelines and batch vs real-time inference.
What does the pipeline look like?
| Stage | Responsibility | Key concern |
|---|---|---|
| Source topic | Raw events | Schema, partitioning key |
| Router | Rules first; decide which events need a model | Keep inference volume proportional |
| Batcher | Group events per inference call | Latency versus throughput |
| Inference consumer | Call the model with bounded concurrency | Backpressure, idempotency |
| Validator | Check outputs against schema and rules | Reject before producing |
| Enriched topic | Events plus inference results | Versioned schema |
| Dead-letter topic | Failures after bounded retries | Replay tooling |
| Result store | Idempotency and lookup | Keyed by event identity |
The router is the stage most often omitted and the one that decides whether the pipeline is affordable.
Why route before inferring?
Because most events do not need a model. In a typical stream a large share can be handled by deterministic rules: known senders, structured fields, patterns already classified, events below a threshold of interest. Sending those to a model costs money and adds latency for no gain.
A routing stage applies rules first, sends only the remainder to inference, and tags each event with the path it took. That keeps inference volume proportional to value and makes cost predictable. Sending every event to a model is the single most common cost failure in streaming AI. See the model routing and cost control whitepaper.
How is backpressure handled?
By design, not by hoping lag stays small. The inference consumer runs with bounded concurrency matched to the model's throughput and the rate limits of any external provider. When the stream outpaces inference, lag grows, which is acceptable and measurable, rather than memory growing, which is not.
Consumer lag becomes the primary operational metric. Alerts fire on lag growth rate rather than absolute lag, since a steady backlog under control differs from one accelerating. Scaling adds consumers up to the partition count and, beyond that, requires repartitioning or a faster model tier for the routed subset.
How does batching change things?
It changes the economics more than any other single decision. Many enrichment tasks, classification, extraction of a few fields, sentiment, and routing, accept several events in one model call, and batching amortises the prompt overhead across them. Throughput improves substantially, cost per event falls, and the latency increase per event is usually acceptable for enrichment workloads.
Batching interacts with ordering and idempotency: a batch that partially fails needs per-event result handling, not batch-level retry, or successful events get reprocessed alongside failed ones.
How are ordering and idempotency handled?
Kafka delivers at least once by default and reprocesses on consumer rebalance, so every inference result must be safe to produce twice. The pattern that works keys results on event identity, stores them in a result store or uses an idempotent producer with the event key, and checks before inferring whether a result already exists.
Ordering within a partition is preserved by Kafka but not by a consumer that processes concurrently, so where ordering matters, concurrency is per partition rather than per event, or the enriched topic carries sequence information that downstream consumers use to reorder.
What schema governance applies?
The same as any topic, with attention to the model's contribution. Enriched events carry the original fields plus inference outputs: classifications, extracted values, confidence scores, the model and prompt version that produced them, and the routing path. All of that belongs in a versioned schema in the registry, so downstream consumers know what they are reading and a change to the model's output shape is a schema evolution rather than a surprise.
Recording model and prompt version per event is what makes later analysis possible: when quality shifts, the version field says whether a model change caused it. See data contracts for ai.
How are failures handled?
Through dead-letter topics, not consumer-side retry loops. An event that fails inference after a bounded number of attempts goes to a dead-letter topic with the error attached, and the consumer moves on, because retrying indefinitely stalls the partition and everything behind the failing event.
Replay tooling reprocesses dead-lettered events once the cause is fixed, whether a provider outage, a malformed event, or a prompt regression. The dead-letter topic's volume is an operational signal worth alerting on.
How is quality monitored on a stream?
By sampling. Every event cannot be reviewed, but a sampled fraction can be scored against a rubric, automatically and by people, and confidence distributions across the stream can be monitored for drift. A drop in average confidence or a shift in class distribution usually precedes a quality complaint. See what is continuous evaluation.
How is cost controlled?
Through routing, batching, model tier selection for the routed subset, and per-topic cost attribution. Cost per enriched event is the metric, and it should be visible alongside lag and throughput. A pipeline whose cost per event rises without a corresponding quality change has usually lost its routing discipline or its batch size.
What does the build sequence look like?
One week on schemas for source and enriched topics and the result store. One week on the router with rule coverage measured. Two weeks on the batched inference consumer with backpressure, idempotency, and dead-lettering. One week on sampling-based quality monitoring. Then tuning batch size and concurrency against measured lag and cost.
What goes wrong?
Every event sent to a model. Unbounded concurrency that exhausts provider rate limits. Retry loops that stall partitions. Results produced twice on rebalance. Enriched topics without schemas. Model versions not recorded, so quality drift is unexplainable. And lag alerts on absolute values that fire constantly or never.
How FISTA Solutions helps
FISTA Solutions builds Kafka AI pipelines with rule-first routing, batched inference under bounded concurrency, idempotent results keyed by event identity, versioned enriched schemas recording model provenance, dead-letter handling with replay, and cost per event visible from day one, through AI enablement, AI agents, and forward deployed engineers working with platform teams. The record behind the approach is 150+ projects for 50+ companies with 99.9% uptime.
To put inference on a stream without drowning in lag or cost, message FISTA on WhatsApp, or read batch vs real-time inference.
Share-ready article cover
Download the generated social format.
Clear answers
Questions raised by this field note.
Straightforward guidance for evaluating scope, fit, and the next step.
01What is the core design challenge of AI on Kafka?
Speed mismatch. Kafka delivers events faster than any model infers, so a naive consumer falls behind without bound. The design must batch, bound concurrency, apply backpressure, and decide which events need inference at all, or the lag grows until the pipeline is useless.
02How is idempotency handled?
By keying every inference result on the event's identity and making the enrichment operation safe to repeat, because Kafka delivers at least once by default and consumers reprocess on rebalance. A result store or idempotent producer pattern prevents duplicate enrichment and duplicate downstream effects.
03How does batching change the economics?
Substantially. Many classification and extraction tasks accept several events per model call, and batching amortises prompt overhead and improves throughput, often by an order of magnitude, at the cost of slightly higher latency per event, which most streaming enrichment tolerates.
04What should go to the model?
Only events that rules cannot handle. A routing stage applies deterministic logic first, sends the remainder to inference, and records which path each event took. Sending every event to a model is the most common cost failure in streaming AI.
05How are failures handled?
Through dead-letter topics for events that fail inference after bounded retries, with replay tooling to reprocess them once the cause is fixed, rather than blocking the partition or retrying indefinitely in the consumer, which stalls everything behind the failing event.
Continue exploring
Related capabilities
Start with the hard problem
Need the outcome owned, not merely analyzed?
Tell us where delivery is constrained. WeтАЩll map the fastest credible path from intent to verified production.