Contents
Figure 1: Kafka's architecture — topics, partitions, and brokers — forms the transport layer that lets ML models score events in real time
Recall the batch-vs-streaming article's central lesson: Kafka isn't a processing framework, it's transport. This article is about actually understanding that transport layer well enough to plug an ML model into it — the specific, incremental pattern that lets a fraud model, a recommendation engine, or an anomaly detector score real events as they happen, without you rebuilding your entire platform to get there.
Kafka just hit version 4.2.0 in February 2026, and the most consequential recent change is architectural, not feature-driven: Kafka 4.0 fully removed ZooKeeper, replacing it with KRaft mode for metadata management. If you're following an older tutorial that mentions ZooKeeper setup steps, that's genuinely outdated — current Kafka doesn't need it at all.
By the end of this guide, you'll understand Kafka's core concepts well enough to reason about them, build a working producer-consumer pipeline, and know the specific, low-risk pattern for adding ML inference to a stream without a "big bang" migration. IMO, that incremental-adoption framing from Confluent's own guidance is genuinely the most practical advice in this whole space :)
What Kafka Actually Is
Apache Kafka is a distributed event streaming platform — a highly scalable, fault-tolerant, durable commit log where events get published to topics, partitioned for speed, and stored across a cluster of brokers. Originally built at LinkedIn specifically to handle its own internal data pipelines, it's since become the backbone of real-time infrastructure at companies like Netflix, Uber, and Airbnb.
The throughput numbers are genuinely worth internalizing as context: Shopify alone processes over 1.5 million Kafka messages per second during peak traffic events like Black Friday. This isn't a niche tool — it's built specifically for the scale where a database or a simple message queue would fall over.
The Core Concepts You Actually Need
Five ideas, and understanding these genuinely unlocks everything else about how Kafka works.
Topics — a named category events get published to, conceptually similar to a table name, except records are appended, never updated in place.
Partitions — a topic is split into partitions specifically for parallelism and speed; each partition is an ordered, immutable sequence of records.
Brokers — the actual Kafka servers storing data and serving clients; a real cluster runs multiple brokers for redundancy and scalability, not just one.
Producers and Consumers — producers publish events to topics; consumers subscribe to and read from them. Consumers are organized into consumer groups, where each partition gets consumed by exactly one consumer within a group — this is the actual mechanism enabling load balancing across multiple consumer instances.
Offsets — a unique, sequential ID assigned to each record within a partition. Consumers track their own position using offsets, which is precisely what enables exactly-once or at-least-once processing guarantees, depending on how you configure offset commits.
Kafka exposes four key APIs on top of these concepts: the Producer API, Consumer API, Streams API (for actual processing, recall the batch-vs-streaming article's distinction), and Connector API (for plugging in external systems without writing custom producer/consumer code by hand).
Setting Up a Local Kafka Cluster
Since ZooKeeper is gone as of Kafka 4.0, current setup is genuinely simpler than what older tutorials describe. Docker Compose is the standard way to get a reproducible local cluster running.
# docker-compose.yml
services:
kafka:
image: apache/kafka:4.2.0
ports:
- "9092:9092"
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@localhost:9093
docker compose up -d
Notice there's genuinely no separate ZooKeeper service defined here — KAFKA_PROCESS_ROLES: broker,controller means this single node handles both jobs itself via KRaft, exactly the simplification that makes current Kafka setups meaningfully lighter than what pre-4.0 tutorials required.
Writing a Producer
from kafka import KafkaProducer
import json
producer = KafkaProducer(
bootstrap_servers="localhost:9092",
value_serializer=lambda v: json.dumps(v).encode("utf-8")
)
order_event = {
"order_id": "ord_10234",
"amount": 149.99,
"customer_id": "cust_5521",
"timestamp": "2026-09-09T14:22:00Z"
}
producer.send("orders", value=order_event)
producer.flush()
This is genuinely the entire "publish an event" pattern — connect to the cluster, serialize your data, send it to a named topic. Every real Kafka-fed ML pipeline starts with something conceptually identical to this, just with a genuine event source (a payment processor, a web app, a sensor) instead of a hardcoded dictionary.
Writing a Consumer
from kafka import KafkaConsumer
import json
consumer = KafkaConsumer(
"orders",
bootstrap_servers="localhost:9092",
value_deserializer=lambda v: json.loads(v.decode("utf-8")),
group_id="fraud-detection-group",
auto_offset_reset="earliest"
)
for message in consumer:
order = message.value
print(f"Processing order {order['order_id']} for ${order['amount']}")
That group_id parameter is doing genuinely important work — it's what makes this consumer part of a consumer group, meaning if you run multiple instances of this same script, Kafka automatically distributes partitions across them for load balancing rather than every instance processing every message redundantly.
Where ML Actually Plugs In: The Incremental Pattern
Here's the genuinely important guidance from Confluent's own current documentation, and it directly echoes this series' MLOps-arc philosophy: the most effective way to adopt streaming ML is not by rebuilding your entire platform, but by adding a single, high-value inference step to your existing data flow. This lets you transition from batch processing to real-time decisions incrementally, without the risk of a "big bang" migration.
from kafka import KafkaConsumer, KafkaProducer
import json
import joblib
model = joblib.load("fraud_model.pkl") # recall the MLOps/CI-CD articles' trained, versioned model
consumer = KafkaConsumer(
"orders",
bootstrap_servers="localhost:9092",
value_deserializer=lambda v: json.loads(v.decode("utf-8")),
group_id="fraud-scoring-service"
)
producer = KafkaProducer(
bootstrap_servers="localhost:9092",
value_serializer=lambda v: json.dumps(v).encode("utf-8")
)
for message in consumer:
order = message.value
features = [[order["amount"], len(order["customer_id"])]] # simplified
fraud_score = model.predict_proba(features)[0][1]
order["fraud_score"] = fraud_score
producer.send("scored-orders", value=order)
if fraud_score > 0.8:
producer.send("fraud-alerts", value=order)
Notice this consumer-then-producer pattern: read raw events from orders, enrich them with a model's prediction, and publish the enriched result to a new topic (scored-orders) that downstream systems can consume — genuinely the same "specialized tools, each doing one job" philosophy from earlier in this series, just applied to a streaming pipeline instead of a batch one. Anything downstream needing the fraud score subscribes to scored-orders rather than needing to know anything about how it was computed.
Kafka Streams: When You Need Real Processing, Not Just Pass-Through
Recall the batch-vs-streaming article's point that Kafka itself doesn't process anything — Kafka Streams is the library that actually does, letting you build stateful transformations, aggregations, and windowed computations directly against Kafka topics without standing up a separate Spark or Flink cluster.
# Faust — a Python alternative to Java-based Kafka Streams
import faust
app = faust.App("fraud-detector", broker="kafka://localhost:9092")
orders_topic = app.topic("orders", value_type=dict)
@app.agent(orders_topic)
async def process_orders(orders):
async for order in orders:
if order["amount"] > 10000:
print(f"High-value order flagged: {order['order_id']}")
Java is genuinely the performance-guarantee choice for Kafka Streams, while Faust offers a simpler development experience for Python-native teams — worth picking based on your team's existing language comfort rather than assuming Java is mandatory. This directly connects to the batch-vs-streaming article's guidance: Kafka Streams is the right level of complexity for simpler, Kafka-native transformations, with Flink reserved for genuinely complex stateful processing that outgrows what a lighter library comfortably handles.
Kafka Connect: Skip Writing Custom Producers for External Systems
If your data source is already a database, a SaaS API, or a file system rather than something you're writing a custom producer for, Kafka Connect provides pre-built connectors handling that integration without custom code.
Source connectors pull data into Kafka from external systems — a database's change log, an API's webhook stream.
Sink connectors push data out of Kafka into external systems — writing scored predictions into a database, a data warehouse, or a search index.
This directly parallels the ETL pipeline article's extraction step — Kafka Connect is genuinely the streaming-world equivalent of a Python extract_from_api() function, just implemented as a managed, pre-built plugin instead of custom code you maintain yourself.
Exactly-Once Semantics: Why This Matters for ML Specifically
Recall the CI/CD article's quality-gate concern about reproducibility — duplicate processing is a genuinely real risk in streaming systems, and it matters more for ML than you might initially assume.
At-least-once delivery (Kafka's simpler default guarantee) means a message might get processed more than once if a consumer crashes and restarts before committing its offset.
For a fraud-scoring pipeline, processing the same transaction twice could mean double-charging a fraud alert or double-counting a feature aggregation — genuinely consequential, not just a theoretical concern.
Kafka's exactly-once semantics, configured through idempotent producers and transactional consumers, prevent this — directly echoing the ETL pipeline article's idempotency discussion, just implemented at the streaming-transport level instead of the transformation-function level.
Realistic Performance Expectations for ML Inference
Batch inference pipelines running on hourly or daily schedules simply cannot meet the requirements of fraud detection, recommendation engines, dynamic pricing, or anomaly detection systems demanding sub-second response times — this is genuinely the concrete case where the batch-vs-streaming article's decision framework points firmly toward streaming rather than a more frequent batch job.
Feature computation happening inline within the consumer loop needs to stay genuinely fast — a model inference step adding meaningful latency to every single message processed can become the actual bottleneck in an otherwise high-throughput pipeline.
This is precisely why the fraud detection architecture from the batch-vs-streaming article split training (batch, nightly, can afford to be slow) from serving (streaming, needs to be fast) — don't let a slow, complex model force your entire pipeline into batch when only the training step genuinely needs that luxury.
Common Mistakes People Make
Following outdated tutorials referencing ZooKeeper setup. Current Kafka (4.0+) uses KRaft mode natively — ZooKeeper is genuinely gone, not just optional.
Treating Kafka itself as a processing engine. It's transport — you need Kafka Streams, Faust, Spark Structured Streaming, or Flink layered on top to actually transform or score events, exactly the distinction the batch-vs-streaming article made explicit.
Attempting a full platform migration to streaming instead of the incremental single-inference-step pattern. Confluent's own current guidance is explicit here — add one high-value inference step to an existing flow first, rather than rebuilding everything at once.
Ignoring exactly-once semantics for use cases where duplicate processing genuinely matters. A fraud alert or a billing-adjacent feature computed twice is a real, consequential bug, not a theoretical edge case.
Skipping consumer groups and running duplicate, uncoordinated consumer instances. Without a shared group_id, multiple consumer instances process every message redundantly instead of load-balancing across partitions as intended.
Recommended Books
- Designing Event-Driven Systems by Ben Stopford — the definitive guide to thinking in events, covering Kafka's architecture from first principles through to production patterns. Essential reading if you're building any streaming infrastructure, not just ML pipelines.
- Kafka: The Definitive Guide by Gwen Shapira et al. — the comprehensive reference for Kafka internals, consumer groups, exactly-once semantics, and operational best practices. The chapter on exactly-once delivery alone is worth the price for ML engineers working with streaming.
- Machine Learning Engineering by Andriy Burkov — covers the MLOps lifecycle from training to serving, including the architectural decisions around batch vs streaming inference that this article's incremental pattern addresses.
Wrapping This Up
Kafka gives ML pipelines a durable, high-throughput transport layer for events — topics, partitions, brokers, producers, and consumers form the vocabulary, and the actual ML value comes from a consumer that scores incoming events with a trained model and republishes the enriched result, not from Kafka itself doing any processing. The incremental adoption pattern — one high-value inference step bolted onto an existing flow, not a platform rewrite — is genuinely the practical path most real teams should follow.
Remember that Kafka 4.0+ has fully removed ZooKeeper in favor of KRaft mode, so any setup instructions still describing a separate ZooKeeper service are outdated. FYI, this genuinely closes the loop with the batch-vs-streaming article from earlier in this series — you now have the concrete transport layer and the actual code pattern for the "streaming serving" half of that batch-train/streaming-serve architecture, rather than just the conceptual argument for when to use it :)
Now go take the fraud model concept from this article and actually swap in a real trained model from earlier in this series — even the CartPole DQN's underlying predict() call structurally fits this same consumer-loop pattern. That's genuinely the fastest way to feel how any trained model, regardless of what it does, plugs into a live event stream the same basic way.