Technology guide What is Apache Kafka®?
Apache Kafka® is an open-source distributed event streaming platform. Applications publish events to Kafka and other applications read them at their own pace. Kafka stores those events so several applications can process the same data and return to earlier records when needed.
Apache Kafka use cases
Kafka gives applications a shared stream of events that they can read independently and return to when they need to repeat some processing. Its uses range from coordinating order updates across services to moving database changes into reporting systems.
| Use case | Application requirement | Why Kafka fits |
|---|---|---|
| Order processing | Share order updates with fulfilment and reporting applications | Independent consumer groups process the same retained events |
| Change data capture | Update search indexes and analytical stores from database changes | Carries captured changes to multiple downstream systems |
| Fraud detection and live analytics | Analyse activity as new events arrive | Feeds processing applications that maintain running calculations |
| Telemetry and log pipelines | Keep ingesting during temporary slowdowns in downstream processing | Retains a backlog that consumers can read as they catch up |
Order processing across services
After a customer places an order the fulfilment service can arrange delivery while a separate reporting application updates its sales figures from the same event. Each application uses its own consumer group and can resume reading retained records after a restart. A reporting failure need not stop fulfilment from progressing. Publishing the event still needs to be coordinated with saving the order and consumers must handle retries without creating duplicate shipments. The data integration guide describes how one stream can feed several applications.
Change data capture and data pipelines
A retailer might need product changes in its operational database to appear in its search index and analytical warehouse. A change data capture (CDC) connector can publish supported database changes to Kafka so separate consumers update those destinations without each system repeatedly querying the source tables. Kafka carries the captured changes while the connector supplies the database-specific capture logic. Initial snapshots and schema changes need planning alongside how consumers handle deletions and repeated records.
Fraud detection and live analytics
A fraud-analysis application might count payment attempts for an account over five minutes and compare each new attempt with the recent activity. Kafka supplies the event stream while a processing application built with Kafka Streams maintains the calculation and emits results. Similar processing can update sales totals as purchases arrive. Detection rules belong to the application and need to account for delayed events and repeated records when deciding what to flag.
Telemetry and log collection
A factory can publish machine readings into Kafka for both an alerting application and a historical data store. If the storage consumer slows down then Kafka can retain the incoming readings while that consumer catches up as long as the brokers have sufficient capacity and retention covers the delay. The alerting application can continue reading independently. Our Kafka-to-Cassandra example shows how a sink connector writes events into a database that applications can query.
A direct request or a simpler queue may be easier to operate for a small workflow that does not need independent readers or replay. Choosing Kafka also means planning the available storage and how applications recover after falling behind.
An event example
For an ordering service an event can describe a newly created order with the order identifier as its key. Here we have encoded the value as JSON for readability while Kafka itself stores the key and value as bytes.
Key: order-1042
Value:
{
"order_id": "order-1042",
"status": "created",
"total": 49.95,
"currency": "GBP"
} The ordering application is a producer because it publishes the record. A fulfilment application is a consumer because it reads it. The named stream they use is a topic called orders.
Records also carry timestamps and can include headers for information such as a correlation identifier. Applications agree on how to encode and interpret the data and can use a schema to check field changes before they reach consumers.
Topics and partitions
A topic is divided into partitions. Each partition is an ordered sequence of records and every record has an offset identifying its position within that partition.
order-1042createdorder-1042paidorder-1043createdorder-1043cancelledThe created record for order-1042 comes before its paid record in partition 0. There is no single ordering across both partitions so comparing their offset values does not tell us which order was created first.
A producer using consistent key-based partitioning can keep records for the same order together. The actual partition depends on the client and its partitioning configuration. Changing the partition count can change where a key is sent. The topics guide explains those choices.
Consumer groups and replay
The fulfilment service can run two consumer instances in one consumer group. For the conventional consumer-group model each partition is assigned to one member at a time. Kafka can assign partition 0 to one instance and partition 1 to the other so they share the work.
| Group | Consumer instance | Assigned partitions |
|---|---|---|
| fulfilment | Worker A | Partition 0 |
| fulfilment | Worker B | Partition 1 |
| sales-reporting | Reporting worker | Partitions 0 and 1 |
The reporting worker belongs to a different group and reads both partitions independently. Its progress does not remove records that fulfilment still needs. Each group commits offsets to record where it should resume.
Adding a third fulfilment consumer would leave one idle in this two-partition example. More consumer instances only add useful parallelism when there are partitions available to assign to them. Our consumer groups guide covers assignment and rebalancing.
A consumer can return to an earlier offset while the records are still retained. That is useful for rebuilding a report after fixing its processing logic. Replaying a payment event needs more care because the application must avoid charging the customer again.
Committing progress after processing can cause a record to be processed again if the consumer fails before the commit. Kafka transactions can coordinate Kafka writes and consumer offsets for exactly-once processing within supported workflows. They do not automatically make an external payment API exactly once. The exactly-once guide explains those limits.
Brokers and KRaft controllers
A Kafka broker stores partition replicas and handles requests from producers and consumers. Each partition has a leader replica that accepts writes while follower replicas copy its log. Different partitions can have their leaders on different brokers.
KRaft controllers manage cluster metadata and coordinate changes such as partition leader elections. A controller quorum elects an active controller and replicates the metadata log. Application records remain on the brokers rather than passing through the controller as a central message relay.
| Component | Responsibility |
|---|---|
| Producer | Sends records to the leader for the chosen partition. |
| Broker | Stores partition replicas and serves produce and fetch requests. |
| KRaft controller quorum | Maintains cluster metadata and coordinates leadership changes. |
| Consumer | Fetches records and tracks its processing progress. |
Kafka 4.x uses KRaft without ZooKeeper and can combine broker and controller roles in one process for development deployments. Dedicated controllers separate those responsibilities in larger deployments. The KRaft architecture guide explains the metadata quorum.
Replication and broker failures
A partition with a replication factor of three has three copies on different brokers. The in-sync replica set (ISR) contains the leader and followers that are sufficiently caught up. Producers using acks=all wait for the required replication acknowledgement rather than just the leader's local append.
| Partition | Broker 1 / Zone A | Broker 2 / Zone B | Broker 3 / Zone C |
|---|---|---|---|
| orders-0 | Leader | Follower | Follower |
| orders-1 | Follower | Leader | Follower |
Consider replication factor 3 with min.insync.replicas=2 and producer acks=all. If broker 1 fails while both followers are in sync then partition 0 can elect a replacement leader. Writes can resume with two in-sync replicas provided the controller quorum remains available.
If only one in-sync replica remains then writes with acks=all are rejected because the configured minimum cannot be met. The precise acknowledgement behaviour also depends on the Kafka version and configuration. The replication guide covers these settings.
Putting three copies in one availability zone still leaves them exposed to that zone failing. Our Kafka multi-AZ durability guide follows replica placement and one-zone failure scenarios.
Retention and log compaction
Reading a record does not delete it. Kafka stores partition logs in segment files and applies the topic's retention policy independently of consumer progress. Time-based or size-based retention can remove old segments even if a consumer has fallen behind.
Log compaction provides another policy that eventually removes superseded values for a key while retaining its latest value. For a topic representing current order state this can avoid retaining every historical update. Compaction is asynchronous and is not a promise that a consumer will see only one record for each key.
The Kafka storage guide covers segment files and cleanup policies. Retention needs enough room for the recovery window you expect consumers to require.
Kafka Connect and stream processing
Kafka Connect runs connectors that move data between Kafka and external systems. A source connector brings data into Kafka and a sink connector writes it elsewhere. Our Kafka-to-Cassandra example shows records arriving in Cassandra tables through a sink connector.
Kafka Streams is a library for applications that process records and maintain state. An application might calculate sales totals in time windows or join an order stream with customer information. Its processing runs in the application rather than as user code inside the brokers.
A schema registry manages schemas used by producers and consumers and can enforce compatibility rules for changes. It is a separate service from the Kafka brokers. Our schema registry comparison includes a worked schema evolution example.
Apache Cassandra can serve application queries over data delivered through Kafka. Choosing how to store that data still depends on the queries the application needs to run.
Operating Kafka with AxonOps
A running broker does not tell you whether an order reached fulfilment. You also need to follow consumer progress and check replication health alongside the connectors that move data into other systems.
AxonOps combines Kafka monitoring with topic and consumer-group visibility. It also provides Connect management and schema registry tools so you can follow problems across the services involved.
- Understand consumer lag and how to measure whether an application is catching up.
- Monitor and recover Kafka Connect tasks with AxonOps using connector status and task-level information.
- Configure Kafka security with TLS and authenticated access backed by ACLs.
- Choose Kafka monitoring metrics for broker requests and replication alongside application health.