Start with an order
An online store accepts order O-1042. Billing, inventory, customer notifications and analytics all need to react. Calling all four synchronously makes checkout depend on all four services. Writing an OrderPlaced event to a durable stream lets each service catch up independently. The event reports a fact that already happened; a command such as ChargePayment asks another service to do something.
Kafka stores a sequence of records that applications can read again. A consumer does not delete a record by reading it. This makes replay and independent subscribers possible, but also means the application must manage duplicate processing, old events and changes to event structure. Kafka is neither a relational database nor a guarantee that every downstream operation succeeds.
Experiment before reading
The lab starts with six retained order events. Green records are processed by the selected group. Position means the next record that group will fetch; commit means its durable restart checkpoint. Click Poll / process, then Crash / restart, then poll again before committing.
The vocabulary and architecture
| Term | What it means in the store |
|---|---|
| Record / event | A key, serialized value, timestamp and optional headers; example: customer A placed order O-1042 |
| Topic | A named event stream, such as orders |
| Partition | One ordered append-only log within the topic |
| Offset | A record's position inside one partition, not a topic-wide sequence number |
| Broker | A server storing partition replicas and serving clients |
| Leader | The replica coordinating a partition's normal writes |
| Follower | A replica fetching the leader's data |
| Producer | A client serializing and sending records |
| Consumer group | A logical subscriber with partition assignments and its own committed offsets |
| Controller | A server participating in management of cluster metadata |
Checkout → producer → orders topic
P0: [B,0] [B,1] ... → billing C1
P1: [C,0] [C,1] ... → billing C2
P2: [A,0] [A,1] ... → billing C1
↘ analytics group (own offsets)
Each partition has replicas on multiple brokers.
Controllers manage metadata; brokers carry the event data.
In modern Kafka 4.x, KRaft maintains cluster metadata through a controller quorum. It replaces the ZooKeeper-based architecture found in many older tutorials. Metadata consensus and partition data replication are different mechanisms. A combined broker/controller is convenient for local learning; dedicated roles and appropriate quorum placement deserve production planning. Losing a data broker and losing the controller quorum have different consequences.
Follow one event end to end
- Checkout creates a stable event ID and a business key, normally customer ID or order ID. It serializes an agreed schema; Kafka does not interpret the business payload.
- The producer discovers cluster metadata and the destination partition's leader. Bootstrap servers are discovery entry points, not a single permanent routing proxy.
- The partitioner chooses a partition. A stable key with unchanged partitioning configuration normally preserves that key's routing. Null-key behavior depends on the client and partitioner. Increasing partition count can change key placement; plan migrations when per-key order matters.
- The producer buffers, batches and optionally compresses records. It sends a batch to the partition leader. Batching trades waiting time against network efficiency; a timeout or retry can make the application's result ambiguous.
- The leader appends records, and followers replicate them. Acknowledgment depends on the configured durability requirements; replicated does not mean processed by billing.
- Billing fetches records, validates them and performs a business operation. Polling advances a consumer's local position before the application necessarily completes its work.
- After successful processing, billing commits the next offset. Finishing offset 7 means the restart checkpoint is 8. A second group can still read offset 7 independently.
- Retention or compaction later cleans the log. Offsets can have gaps; do not rely on contiguous numbers for business counts.
Partitions, ordering and parallelism
A topic with three partitions provides three independent ordered sequences. It does not provide a total order across the topic. Events at P0:10 and P1:10 are unrelated positions. Choose the key around the entity whose order you need, rather than hoping timestamps can reconstruct one global sequence.
Within a conventional consumer group, one partition is assigned to one member at a time. With three partitions and two members, a member can own two partitions. With four members, at least one is idle. A different group can read all three partitions again. Assignments change during a rebalance; production behavior depends on the protocol and assignor. The lab deliberately uses simple round-robin ownership.
A hot customer key can overload one partition even when the cluster is mostly idle. Salting the key can distribute work but sacrifices straightforward ordering for that entity. Adding partitions also affects routing, metadata and resource use. Start with expected throughput, key distribution, consumer parallelism and retention requirements rather than “more is always faster.”
Exercise: In the lab, publish A three times. Which partition changes? Switch groups, add a fourth member and explain why it does not create a fourth ordered log.
Replication, acknowledgments and failure
A replication factor of three means three copies of each partition, usually placed on separate brokers. The in-sync replica set (ISR) contains replicas sufficiently caught up according to Kafka's rules. A slow follower can leave it. A returning broker must catch up; the lab makes this instantaneous for clarity.
| Setting | Meaning | Risk to reason about |
|---|---|---|
| acks=0 | Producer does not wait for a broker acknowledgment | Cannot reliably establish whether a write succeeded |
| acks=1 | Leader acknowledges its append | Leader failure before replication can lose acknowledged data |
| acks=all | Wait for the current ISR acknowledgment requirement | Availability depends on ISR and minimum replica settings |
| min.insync.replicas=2 with RF=3 | Require enough in-sync copies for these durable writes | A depleted ISR causes writes to fail rather than silently weakening the requirement |
| enable.idempotence=true | Deduplicate qualifying producer retry sequences | Does not deduplicate a new business operation or an external database write |
acks=all is not “wait for all replicas forever,” and a successful acknowledgment is not proof of completed payment. Durability depends on the replica policy and failure model. Avoid electing an out-of-date replica merely to regain availability without understanding possible data loss; leader-election behavior also depends on the Kafka release and enabled features.
Lab: stop one broker, publish, then stop another and publish again. Restore a broker. This lab keeps acks=all, RF=3 and minimum ISR=2 fixed and models only clean takeover.
Offsets, lag, restart and replay
An offset commit is a checkpoint, not a business transaction. Kafka stores consumer-group checkpoints in its internal offsets topic. A member can have fetched farther than it has processed, and processed farther than it has committed. Treat those as three separate positions when debugging.
Committed lag is often approximated as log end minus committed offset for each partition. It is a distance in offsets, not inherently seconds or an exact count of business events. Measure event age and processing delay as well. The lab uses dense records without compaction, so its displayed lag is easy to interpret.
auto.offset.reset applies when no valid committed offset exists or it is out of range; it does not force every restart to begin at the start. “Earliest” means the earliest retained record, which may no longer be offset zero. A replay can repeat side effects and may encounter old schemas. Stop the group for an administrative reset, choose the range, validate the new consumer and protect downstream systems from the extra load.
Long processing can exceed the polling interval and cause reassignment. If multiple workers process a partition concurrently, commit only through a completed contiguous prefix; checkpointing beyond unfinished work can skip it on restart. On partition revocation, stop admitting work and handle outstanding tasks according to the client protocol.
Delivery guarantees: where duplicates come from
| Processing order | Failure window | Result |
|---|---|---|
| Commit, then execute side effect | Crash after commit, before execution | Work can be skipped: at-most-once behavior |
| Execute side effect, then commit | Crash after execution, before commit | Work can repeat: at-least-once behavior |
| Transactionally write Kafka output and consumed offsets | Abort or retry the Kafka transaction | Kafka read-process-write can avoid duplicate visible output with the correct isolation |
Kafka transactions can couple Kafka output records and consumed offsets; consumers using read_committed avoid reading aborted transactional records. They do not automatically include an arbitrary SQL transaction, an email or a payment-provider request. Producer idempotence and application idempotency solve different problems.
For a SQL sink, one useful application design is a unique event ID recorded in the same database transaction as the side effect:
BEGIN database transaction
insert event_id into processed_events (unique key)
if already present: do not repeat the business update
otherwise: apply the business update
COMMIT database transaction
commit the Kafka offset after the database transaction succeeds
A crash between the database commit and offset commit causes redelivery, but the unique event ID suppresses the repeat update. This is an application design, not a universal broker promise. Payment calls need the provider's idempotency mechanism too. Check identifier scope, retention of deduplication records and replay behavior.
Retention, compaction and schemas
Retention can remove old segments by time or size. Consumption does not trigger deletion. If a consumer falls behind the retained history, replaying missing events is impossible without another source or archive.
Compaction eventually preserves the latest surviving value per key instead of every update forever. A null-valued tombstone represents deletion for a key. Compaction is asynchronous; older versions can coexist until cleaned, and tombstones have their own retention behavior. A compacted log can help reconstruct current state but is not a permanent full audit history.
Design an event envelope deliberately:
{
"event_id": "evt-O1042-v1",
"event_type": "OrderPlaced",
"schema_version": 1,
"occurred_at": "2026-10-03T10:00:00Z",
"order_id": "O-1042",
"customer_id": "A",
"amount_minor": 249900,
"currency": "INR"
}
This is our example schema, not Kafka's required format. Prefer integer minor units over floating-point money. Distinguish event time from ingestion time. Use schema compatibility tests before rollout: an optional added field can be safer than changing an existing field's meaning. Avro, Protobuf and JSON Schema are choices; registry products and serializers are ecosystem components, not a mandatory feature of every Kafka installation. Never casually put credentials or unnecessary personal data in a long-retained topic.
Connect, CDC, outbox and processing
Kafka Connect runs source and sink connectors for moving data between Kafka and external systems. Connector tasks, their offsets, converters and operational errors still need management. Installing a connector does not turn an incompatible database schema into a correct event contract.
The dual-write trap is simple: the application commits an order in SQL and crashes before publishing its event, or publishes and then rolls back SQL. A transactional outbox records the order and its pending event in one database transaction. A CDC connector or publisher moves outbox records to Kafka. Downstream consumers still need duplicate protection; an outbox is not an excuse to assume one delivery.
Kafka Streams is an application library for transforming Kafka records, maintaining local state, joining streams and tables, and producing results. Repartitioning and changelog topics support the computation. Apache Flink is a separate processing engine; it is not a substitute for a broker's durable transport. Event-time windows group by when events happened. Watermarks estimate progress so a processor can decide how to treat late arrivals; they do not guarantee that a late event cannot appear.
Order DB + outbox → CDC / Connect → orders
├─ billing → ledger DB
├─ notifications → email provider
└─ stream processor → revenue-per-minute
→ dashboard sink
Retries need a policy: exponential backoff for transient failures, bounded attempts, and a dead-letter path for poison events. Moving failed records to retry topics can change ordering. An audit log and an owned replay procedure matter more than quietly collecting an unmonitored DLQ.
Batch, streaming and micro-batch processing
A batch job processes a bounded input, such as yesterday's transactions, and finishes. A streaming job handles an ongoing input and maintains state as records arrive. The choice is about input and result freshness, rather than a particular product name.
| Approach | Useful example | Main design question |
|---|---|---|
| Batch | Reconciled daily finance report | Can the business wait for the next completed run? |
| Streaming | Fraud signals during payment | How are state, late events and failures handled? |
| Micro-batch | Dashboard refreshed every short interval | Is the interval acceptable, including processing delay? |
Continuous processing requires operating the pipeline continuously, not merely making it fast once. Batch can simplify recovery through repeatable input snapshots. Micro-batches trade some freshness for grouping work; they still need checkpoints and duplicate-safe outputs. A common design combines immediate operational signals with later batch reconciliation, with an explicit rule for correcting earlier results.
Build a local pipeline
This exercise uses the official Kafka 4.2.0 Docker image as a pinned learning baseline, not a claim that it is the newest release. Docker must be installed. Use a local-only broker; this one-node setup has no production redundancy. Commands are reference instructions; the website lab does not execute Docker.
docker run --rm --name kafka-lesson -p 127.0.0.1:9092:9092 apache/kafka:4.2.0
Wait for startup. In another terminal, create and describe the topic:
docker exec kafka-lesson /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 --create --topic orders \
--partitions 3 --replication-factor 1
docker exec kafka-lesson /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 --describe --topic orders
Expected: three partition descriptions, each with one replica. Do not use minimum ISR=2 on this one-replica exercise.
docker exec -it kafka-lesson /opt/kafka/bin/kafka-console-producer.sh \
--bootstrap-server localhost:9092 --topic orders \
--property parse.key=true --property key.separator=:
Type A:order-1042, B:order-1043, then A:order-1044, each on a new line. Consume and display partition/offset information:
docker exec -it kafka-lesson /opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 --topic orders --group billing \
--from-beginning --property print.key=true \
--property print.partition=true --property print.offset=true
Run the consumer in another terminal with --group analytics. Both groups should see the history. The A events should share a partition with this stable configuration; do not expect the real client to use the lab's toy hash. Arrival order across partitions need not match your typed order. Add more billing consumers and observe redistribution. Inspect group offsets:
docker exec kafka-lesson /opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 --describe --group billing
# After ending the terminal experiments:
docker stop kafka-lesson
The console consumer helps inspect delivery; it is not an implementation of your billing transaction. Its commit behavior differs from the lab's explicit manual controls. Stopping this disposable container removes its unmounted local data.
Other platforms and their differences
Compare transport and storage separately from processing engines and managed hosting. A shared word such as “topic” or “exactly once” does not imply identical semantics. The selection guidance below is our synthesis of the linked product documentation, not a throughput benchmark or price comparison.
| Platform | Model and useful fit | What to check before choosing |
|---|---|---|
| Apache Kafka | Partitioned durable logs; replay, multiple independent groups, extensive integration ecosystem | Partition/key strategy, storage sizing, operations, client compatibility and transaction scope |
| Redpanda | Kafka API-compatible streaming platform with its own implementation | Validate required APIs, clients, features and operational behavior; compatibility is not identical internals |
| Apache Pulsar | Messaging and streaming with brokers separated from BookKeeper storage; multiple subscription patterns | Subscription mode, ordering, operational components and required integration features |
| NATS JetStream | Persistent streams and consumers alongside NATS messaging | Ack policy, replay/retention mode, limits and deduplication scope; distinguish JetStream from transient Core NATS |
| RabbitMQ Streams | Persistent replicated logs with non-destructive consumption and replay | Streams versus classic/quorum queues, native protocol/client behavior and superstream partitioning |
| Amazon Kinesis Data Streams | Managed AWS streams partitioned into shards | Capacity mode, partition-key distribution, consumer model, retention and regional costs |
| Google Cloud Pub/Sub | Managed topics/subscriptions and acknowledgment-based delivery | Ordering-key configuration, replay retention and regional restrictions on optional exactly-once delivery |
| Azure Event Hubs | Managed partitioned event ingestion with independent consumer groups and Kafka endpoint support | Supported protocol features, capacity tier, retention, checkpoint ownership and ecosystem integration |
| Redis Streams | Stream entries and consumer groups inside Redis | Pending-entry recovery, trimming, persistence/replication settings and memory/operational constraints |
Managed Kafka offerings such as Amazon MSK, Confluent Cloud and hosted Kafka providers move some infrastructure work to the provider. They do not remove schema design, consumer correctness or application monitoring. Flink and Spark Structured Streaming process streams; Debezium captures database changes. They can sit alongside a broker rather than competing with it for the same job.
Architectural families in more detail
Kafka-compatible platforms: Kafka and Redpanda
Kafka applications normally rely on a partitioned log, consumer groups and the Kafka wire protocol. Redpanda implements that protocol with a different server architecture. Existing client code can therefore be a starting point, but migration still requires testing the features your application actually uses: transactions, admin APIs, connectors, authentication and operational tooling. Keeping the same producer API is different from keeping the same cluster runbook.
A useful migration trial replays representative data into a test cluster, validates ordering and retry behavior, and compares end-to-end latency with identical payloads and durability requirements. This is our suggested evaluation procedure, not a claim that one implementation is universally faster.
Pulsar: subscriptions and separate storage
Pulsar separates brokers serving clients from the BookKeeper layer storing persistent data. Its subscription choices change delivery behavior: Exclusive gives one consumer access, Failover has an active consumer with standbys, Shared distributes messages among consumers, and Key_Shared routes by key under its specific rules. Shared delivery does not give the same ordering expectation as one ordered partition consumed serially. Retaining already acknowledged messages for later replay requires a suitable retention policy; Kafka's retention assumptions should not simply be carried over unchanged.
Choose the subscription mode from the work's needs, then assess the operational cost of brokers, storage and metadata components. “Separate compute and storage” is an architectural property, not proof that operating the deployment will be easier for your team.
NATS JetStream: subjects, retained streams and acknowledgments
Core NATS messaging and JetStream persistence serve different needs. A JetStream stream captures messages on configured subjects, and a consumer defines how stored messages are delivered and acknowledged. Retention limits and delivery policy determine what can be replayed; acknowledgment policy determines when unfinished messages can be sent again.
For a worker fleet, our suggested review questions are: can workers fetch at their own pace, how is a crashed worker's message recovered, what happens when the stream reaches its limit, and who owns the retry policy? Do not assume that an acknowledgment proves the worker's external database transaction succeeded. Subject routing is also a different abstraction from Kafka's key-to-partition assignment.
RabbitMQ: distinguish streams from queues
RabbitMQ can serve traditional work-queue designs as well as stream use cases. A RabbitMQ stream retains a log with non-destructive reads; consumers can attach at offsets or timestamps. Superstreams distribute a logical stream across partitions. These are distinct from treating an acknowledged queue message as completed work.
Our practical selection advice is to start with the workload: competing background jobs need a work distribution and retry story; multiple independent readers of retained history need a stream/replay story. Choosing RabbitMQ does not require choosing the same data structure for every service. Test the chosen client and protocol, including recovery and offset tracking, rather than assuming a Kafka consumer example maps directly onto RabbitMQ.
Kinesis: an AWS-managed shard model
Kinesis Data Streams routes records into shards using partition keys. Consumers track progress through the shard data, and Kinesis Client Library applications coordinate shard processing with checkpoints. Capacity mode and resharding behavior matter when traffic grows or becomes skewed.
Our suggested evaluation includes a burst with one disproportionately hot key, a consumer outage followed by catch-up, and a change in shard topology. Calculate cost using the selected capacity mode, traffic, consumer reads and retention; do not compare only an advertised ingress number. AWS operation of the service does not replace application-level idempotency or ownership of a failed downstream sink.
Pub/Sub: subscriptions rather than explicit broker ownership
Google Cloud Pub/Sub distributes topic messages to subscriptions and uses acknowledgment-based delivery. Publishers normally receive at-least-once delivery and best-effort ordering; ordering requires the appropriate configuration. Optional exactly-once delivery is subject to documented scope, including supported pull subscriptions and regional constraints, and does not eliminate the possibility of publishing the same business event more than once.
Our design questions are acknowledgment deadline handling, slow subscriber behavior, retention/replay configuration and the distinction between a message acknowledgment and a business side effect. Avoid importing Kafka's explicit partition-assignment mental model into an API where the service manages distribution differently.
Event Hubs: managed event ingestion with multiple interfaces
Azure Event Hubs is a managed partitioned ingestion service. Consumer groups let applications view the event stream independently, and checkpointing lets an application track progress. Its Kafka endpoint enables supported Kafka-protocol clients, but the deployment tier and documented feature support still matter.
For an Azure analytics pipeline, our suggested assessment is to map producers, reader groups, checkpoint storage, retention and downstream processors before selecting capacity. Test your client configuration and required APIs against the actual service tier; a Kafka-compatible connection string is not evidence that every self-managed Kafka operational capability exists unchanged.
Redis Streams: a log inside an existing Redis deployment
Redis Streams uses stream entry IDs rather than Kafka's partition offsets. Consumer groups track deliveries, and a pending-entries list records work delivered but not yet acknowledged. XACK acknowledges group processing; recovery of pending work and retention/trimming require explicit handling. A trimmed entry's payload can disappear even when bookkeeping still refers to it, depending on the chosen commands and options.
Our suggested fit evaluation starts with an existing Redis deployment and a bounded workload, then checks persistence, failover, memory, trimming and recovery. Redis being familiar does not automatically make it the right long-retention event archive. Validate restart and outage behavior with the actual persistence configuration.
How to choose
- A few background jobs: consider a work queue or database-backed job system first. A replayable distributed log may add more complexity than the workload warrants.
- Many independent subscribers need history: investigate durable streaming logs, key distribution and consumer catch-up requirements.
- Kafka integrations are already mandatory: evaluate Kafka or a compatible platform against the precise client/protocol features you need.
- Cloud-native managed ingestion: compare the cloud's streaming service with managed Kafka using retention, fan-out, quotas, networking and cost estimates.
- Event-time joins or aggregations: choose a processor as well as a transport; checkpointing and late-event policy matter.
Run a workload-specific proof of concept. Measure realistic payloads, skewed keys, consumer downtime, replay, schema changes and failure recovery. A vendor benchmark under different conditions does not establish your application's performance.
Production operation and troubleshooting
A useful first storage estimate is:
ingress bytes/second × retention seconds × replication factor
Example: 1,000 events/second × 1 KB × 86,400 seconds × 3 copies is roughly 259 GB/day before compression and other overhead. This is an illustrative decimal estimate, not a provisioned capacity recommendation. Include indexes, extra topics, compaction headroom, operational slack and peak traffic. A consumer processing 1,200 events/second after a 1,000/second producer resumes closes a 120,000-event backlog in approximately 600 seconds if rates remain stable and other bottlenecks do not intervene.
| Symptom | Investigate |
|---|---|
| Committed lag keeps increasing | Processing throughput, downstream latency, skew, failures and commit frequency |
| Lag is low but the dashboard is late | Event age, batch/window delay, sink latency and premature commits |
| Frequent rebalances | Poll interval, member failures, scaling churn and long blocking work |
| Writes reject or time out | ISR, broker availability, leader movement, network, authentication and quotas |
| One partition dominates | Key skew and uneven workload, not only total cluster capacity |
| Disk pressure | Retention, compression, replication, stuck cleanup and capacity headroom |
| Duplicate business effects | Commit timing, retries and missing sink idempotency |
| Missing history after reset | Retention boundaries, incorrect offsets and topic lifecycle |
Track consumer lag and age, throughput, error rates, request latency, offline/under-replicated partitions, ISR changes, disk/network pressure and controller health. Alert on sustained behavior and assign owners.
Encrypt traffic, authenticate clients, authorize only the topics/groups they need and isolate broker listeners. Separate production environments and quotas. Do not expose an unauthenticated broker publicly. Test recovery with the actual configured version and deployment; snapshots, cross-cluster replication and regional recovery need deliberate plans.
Interview checks and a project
- Does Kafka order every event globally? No. Explain per-partition ordering and the chosen key.
- Why is a consumer idle? More group members than assignable partitions, or an assignment/configuration issue.
- Does committing offset 8 mean event 8 finished? The checkpoint normally means the next offset to consume; after processing 7, commit 8.
- Does idempotent production prevent duplicate payments? No. It covers specific producer retry behavior; protect business side effects separately.
- Can you replay forever? Only retained history remains available.
- Why use an outbox? To avoid inconsistent separate database/event writes; downstream duplication remains a concern.
- What is a watermark? A progress estimate for event-time handling, not a guarantee about all future arrivals.
- Would you install Kafka for sending ten emails a day? Explain the simpler alternatives and the requirements that would justify streaming.
Build an order-event project with an outbox, an idempotent SQL consumer and a revenue dashboard. Demonstrate two independent groups, partition skew, a crash before commit, a replay after schema evolution and bounded poison-event handling. Include a runbook and measurements. Do not claim exactly-once email or payment delivery merely because a Kafka transaction is enabled.
Suggested learning sequence: run the browser experiments, reproduce the local CLI exercise, implement an idempotent sink, then add CDC and windowed analytics. Return to Message Queues for the broader messaging foundation or open the standalone Kafka lab.