February 6, 2026
Building Real-Time Data Pipelines with Apache Kafka
How to structure Kafka-based pipelines for observability, warehouse delivery, and downstream agentic reporting: partition design, delivery semantics, replay, and sink patterns.
Article focus
Real-time pipelines only create value when they are observable, replayable, and designed for the downstream decisions they serve.
Section guide
Real-time pipelines are easy to pitch and hard to operate. The technical challenge is not publishing messages into Kafka. The challenge is building a system that stays observable, replayable, and financially sane once traffic grows.
Start with the downstream decision
Before you create topics or write consumers, define the decision that needs fresh data.
Examples:
- risk scoring that needs minute-level fraud signals
- inventory monitoring that needs rapid anomaly alerts
- support workflows that need agent context from live systems
If no decision becomes materially better with faster data, a batch pipeline is often the better answer. This is the question that decides whether the rest of this article applies to you at all, and it is worth being honest about. A dashboard that a human looks at twice a day does not become more useful because the numbers behind it are four seconds old instead of four hours old.
Design the contract before the code
One of the biggest streaming mistakes is shipping events first and governing them later. In practice, event contracts should be designed before scale arrives.
That means agreeing on:
- event naming
- payload ownership
- versioning rules
- dead-letter behavior
- replay expectations
The goal is not bureaucracy. The goal is making sure downstream teams can trust what the stream means.
Make the schema a build-time concern
A schema registry turns "we agreed on the payload" into something the build can enforce. The part worth understanding is the compatibility mode, because it decides who has to deploy first:
- BACKWARD (the common default): a new schema can read data written with the old one. Consumers upgrade first. Safe for adding optional fields, removing fields.
- FORWARD: the old schema can read data written with the new one. Producers upgrade first.
- FULL: both directions hold. The most restrictive, and the right choice for topics with many independent consumers you do not control.
Picking a mode is not a formality. It is the difference between a field rename that rolls out quietly and one that takes down a consumer at 2am because nobody knew which side had to ship first.
Partition design is the decision you cannot undo cheaply
Partitions are where most Kafka regret accumulates, because the two properties they control pull in opposite directions.
Partition count sets your parallelism ceiling. Within a consumer group, one partition is consumed by exactly one consumer. Ten partitions means at most ten useful consumers; the eleventh sits idle. Size from measured throughput rather than a round number: take the peak events per second you need to absorb, divide by what one consumer sustains against your real workload, and add headroom for growth and for recovery after an outage, when you need to catch up faster than real time.
Over-partitioning is not free either. More partitions mean more open file handles, more metadata, and longer rebalances, which is exactly what you do not want during an incident.
The partition key sets your ordering guarantee. Kafka orders messages within a partition, not across a topic. Keying by account_id means every event for one account arrives in order, which is usually the guarantee that actually matters. Keying by something high-cardinality and meaningless, or leaving the key null for round-robin, gives you throughput and no ordering.
Watch for hot keys. If one tenant produces thirty percent of your traffic and you key by tenant, one partition carries thirty percent of the load and one consumer becomes your bottleneck. A composite key such as tenant_id:entity_id spreads the load while preserving the ordering that matters, as long as per-entity ordering is genuinely enough.
Choose your delivery semantics deliberately
"Exactly-once" is the most misread phrase in streaming. It is worth separating two questions.
Inside Kafka, exactly-once is real. An idempotent producer deduplicates retries, and transactions make a read-process-write cycle atomic:
# producer
acks=all
enable.idempotence=true
max.in.flight.requests.per.connection=5
retries=2147483647
delivery.timeout.ms=120000
compression.type=zstd
linger.ms=20
batch.size=131072
# consumer
enable.auto.commit=false
isolation.level=read_committed
auto.offset.reset=earliest
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor
Two of those lines carry more weight than they look. enable.auto.commit=false means you commit offsets after the work succeeded, not on a timer that can acknowledge messages you never processed. CooperativeStickyAssignor makes rebalances incremental rather than stop-the-world, so adding a consumer no longer pauses every consumer in the group.
Outside Kafka, the moment data lands in a warehouse, a search index, or an API call, transactions no longer cover you. The workable pattern is at-least-once delivery into a sink that is idempotent: carry a stable event_id, and write with a merge rather than an insert. Then a replay produces the same rows instead of duplicates, and you never have to reason about whether a batch was applied once or twice.
Plan for replay on day one
Backfills and consumer restarts are not rare edge cases. They are normal operating events, and the time to design for them is before the first one.
Retention is the first lever. A topic you might need to replay from cannot have seven-day retention if your worst-case recovery window is two weeks:
cleanup.policy=delete
retention.ms=1209600000 # 14 days
For topics that represent current state rather than a history of events, compaction is usually the better fit. It keeps the latest value per key indefinitely and discards superseded ones, so a new consumer can rebuild state from the topic without you storing every intermediate change:
cleanup.policy=compact
min.compaction.lag.ms=3600000
The second lever is the sink. Replay is only safe if reprocessing the same events converges on the same result, which comes back to idempotent writes. A merge keyed on the event identifier makes replay boring:
MERGE INTO analytics.orders AS target
USING staging.order_events AS source
ON target.event_id = source.event_id
WHEN NOT MATCHED THEN INSERT *
WHEN MATCHED AND source.ingested_at > target.ingested_at THEN UPDATE SET *;
Treat observability as a product feature
If a consumer slows down or silently drops malformed events, the business impact may be larger than a full outage, because nobody gets paged for data that is quietly wrong.
Measure at least these:
- Consumer lag, the gap between the log end offset and the committed offset, per partition rather than averaged. An average hides the single stuck partition that is actually your problem.
- Message age, the difference between event time and processing time. This is the one to alert on, because it is the number a business owner understands. "Fraud signals are nine minutes stale" means something; "lag is 40,000" does not.
- Dead-letter rate, as a proportion of throughput rather than a raw count.
- Rebalance frequency. Frequent rebalances usually mean
max.poll.interval.msis too tight for how long processing actually takes.
Give poison messages somewhere to go
A single unparseable message should not stall a partition. Route failures to a dead-letter topic with enough context to diagnose them later, and keep the original offset in the headers so you can find the source event:
orders.v1 -> main topic
orders.v1.dlq -> failures, with headers:
x-error-class, x-error-message,
x-source-partition, x-source-offset, x-failed-at
The DLQ is only useful if someone looks at it. An alert on its rate is what turns it from a graveyard into a feedback loop.
Keep enrichment close to business value
Not every transform belongs in Kafka. Put lightweight event normalization near the stream, where it keeps the contract clean: field renaming, type coercion, dropping obviously malformed records. Move heavyweight modeling into systems that are easier to test, version, and backfill, which usually means the warehouse.
The practical test is whether the logic needs history. Anything that requires joining against months of data, or that analysts will want to revise later, is modeling work and belongs where a rebuild is cheap.
Protect warehouse cost
Streaming into the warehouse becomes expensive quickly if every message triggers a write. Three patterns carry most of the savings:
- Micro-batch rather than per-message writes. Buffer on a size or time trigger, whichever comes first, and write once. This is usually the difference between a warehouse bill that scales with events and one that scales with volume.
- Respect the table layout. Partition and cluster sink tables on the columns queries actually filter on, usually event date plus a tenant or entity key.
- Compact small files. High-frequency streaming writes produce many small files, and query engines pay for that on every scan afterwards.
Design human-readable failure paths
When something breaks, operators need an answer to two questions fast:
- what failed
- what action recovers the system safely
That answer should exist before the outage, not after it. In practice that means a short runbook per topic covering where the DLQ is, what the safe replay procedure looks like, and which downstream tables need rebuilding if you do replay. The team that writes this down before the first incident recovers in minutes rather than hours.
Where AI agents fit
Kafka pipelines become more valuable when paired with AI-assisted summaries or anomaly workflows. A useful pattern is:
- stream events into operational storage
- aggregate into monitoring tables
- trigger an agent or assistant only when thresholds matter
That keeps agent activity focused on high-signal events instead of noisy raw traffic, and it keeps the cost predictable. Pointing a model at the raw firehose is both expensive and worse: the signal a reviewer needs is almost always in the aggregate, not the individual event.
A pre-launch checklist
Before a streaming pipeline carries anything a business depends on:
- the downstream decision is written down, and it improves with fresher data
- the schema is registered with a compatibility mode the team has agreed on
- partition key and count come from measured throughput and a stated ordering guarantee
- the sink is idempotent, so replay converges instead of duplicating
- retention covers the worst-case recovery window, not the typical one
- lag alerts fire on message age, per partition
- a dead-letter topic exists and someone owns its rate
- the replay runbook exists and has been tested once, deliberately, before it was needed
The takeaway
Kafka is a great backbone for real-time systems, but only when the pipeline is designed around trust, recovery, and downstream business decisions.
A "working" stream is not enough. A useful streaming platform is one your operators can reason about under pressure, where a replay is a routine procedure rather than an act of courage, and where the cost of freshness is something the team chose on purpose.
Article FAQ
Questions readers usually ask next.
These short answers clarify the practical follow-up questions that often come after the main article.
Need a similar system?
If this article maps to a workflow your team already operates, the next step is usually a scoped review of the system, constraints, and rollout path.
Book your free workflow review here.
Related articles
View allGrok 4.7 vs Fable 5.1 vs GPT-6 Astra: Real Cost per Task
Microsoft Data Days 2026: A Practical Guide for Data Teams
Database Query Plan Regression Review for Production Teams

