Kafka
Kafka is a big log book for messages. Messages are not deleted after they are read, and each reader keeps its own place. Order is kept only inside one Partition, and in each group every Partition has only one reader.
Author: bezzad
The problem: one order, many services
A customer places an order in the online shop. After that, several services must do some work:
- The email service must send the receipt.
- The warehouse service must prepare the items.
- The report service must update today’s sales.
The simple way is for the order service to call all three directly. But this way has three problems:
- Too much coupling. If the email service is broken, placing an order also becomes slow or breaks.
- Hard to change. For every new service, the order service code must change.
- Sudden load. On a sale day, all services must answer all requests at the same moment.
The better way is this: the order service only records one piece of news, “order placed”. Any service that needs it reads it whenever it is ready. Kafka is the place that keeps this news.
The main idea: a log book, not a queue
Think of Kafka as an account book. Each new piece of news is written on a new line, at the end of the book. Nothing in the middle of the book changes.
- A message is not deleted after it is read. A message stays for a set time (Retention). The default time is one week, and you can change it.
- Each reader has its own place. Each reader keeps a number: “I have read up to this line”. We call this number the Offset.
- So readers do not break each other’s work. The email service and the report service both read all messages, each at its own speed.
- You can read again. A new service can start from the beginning of the book and see all the old news.
Topic, Partition and Offset
Each kind of news has a Topic. For example, all order news goes to a topic named orders.
A single book is slow for a large volume. So each topic is split into several parts. We call each part a Partition. Each Partition is a separate book that is kept on one server (and has several copies on other servers).
How does Kafka know which Partition each message goes to?
- Each message has a key (Key). In our shop, we use the order ID as the key.
- Kafka makes a number from the key (Hash). The remainder of this number divided by the number of Partitions is the Partition number.
- So the same key always means the same Partition. All messages of order o-3 go into one Partition, in the same order they were sent.
- If we do not set a key, messages are spread across Partitions, and no order is kept between them.
Many readers for one job: Consumer Group
The email service runs on several pods. We do not want every pod to send its own email to each customer. So we put all the email pods in one Consumer Group.
The group rules are simple:
- Each group reads all messages of the topic. The email group and the warehouse group both see all orders.
- Inside one group, each Partition has only one reader. So no message is processed twice at the same time in one group.
- One reader can have several Partitions.
- Each group has its own Offset. If the warehouse group is behind, the email group does not wait for it.
An important result of the second rule: the largest number of useful readers in a group equals the number of Partitions. If you have six Partitions and ten pods, four pods stay idle.
Live example
Press the buttons. Raise the number of Consumers up to eight and see what happens. Then send the orders of the big customer:
In the first button, the message key is the order ID. In the second button, the key is the customer ID.
You saw two things in this example:
- When the number of Consumers changes, the group is split again. We call this a Rebalance. During this time, reading messages slows down or stops.
- When one key is very busy, one Partition gets hot. All messages of the big customer are on one Partition. So only one Consumer works on them. More Consumers do not help at all. We call this problem a Hot Partition.
Sending messages (Producer)
In .NET we usually use the Confluent.Kafka library. The message key is the order ID:
using Confluent.Kafka;
using System.Text.Json;
var config = new ProducerConfig
{
BootstrapServers = "localhost:9092",
Acks = Acks.All, // wait until all in-sync replicas have the message
EnableIdempotence = true // retries do not write the same message twice
};
using var producer = new ProducerBuilder<string, string>(config).Build();
var order = new OrderPlaced(OrderId: "o-3", CustomerId: "c-42", Total: 250_000);
await producer.ProduceAsync("orders", new Message<string, string>
{
Key = order.OrderId, // same key => same partition
Value = JsonSerializer.Serialize(order)
});
public sealed record OrderPlaced(string OrderId, string CustomerId, decimal Total);
Two important settings in this code:
- The Acks setting set to All. Kafka answers “saved” only when all in-sync copies have the message. So if one server dies, the message is not lost.
- The EnableIdempotence setting. If the network breaks and the Producer sends again, Kafka does not write a duplicate message into the Partition.
Reading messages (Consumer)
After the work, the reader must tell Kafka “I have done it up to here”. We call this committing the Offset. The time of the Commit is very important:
- If you Commit before processing, and then the service dies, the message is lost. We call this at-most-once.
- If you Commit after processing, and the service dies before the Commit, the message comes again. We call this at-least-once.
- So we choose the second case. Losing a message is worse than a repeat. We solve repeats with an Idempotent reader.
using Confluent.Kafka;
var config = new ConsumerConfig
{
BootstrapServers = "localhost:9092",
GroupId = "email-service",
EnableAutoCommit = false, // we commit ourselves, after the work
AutoOffsetReset = AutoOffsetReset.Earliest, // a new group starts from the beginning
PartitionAssignmentStrategy = PartitionAssignmentStrategy.CooperativeSticky
};
using var consumer = new ConsumerBuilder<string, string>(config).Build();
consumer.Subscribe("orders");
while (!ct.IsCancellationRequested)
{
var result = consumer.Consume(ct);
// Idempotent: the handler checks if this order already got an email.
await emailHandler.HandleAsync(result.Message.Key, result.Message.Value, ct);
consumer.Commit(result); // only after the work is done
}
A few notes about this code:
- The CooperativeSticky setting. In a Rebalance, only the Partitions that need to move are moved, not all of them. So the pause is shorter.
- One Consumer per thread. The Consumer object in Confluent.Kafka is not built to be used by several threads at the same time.
- Idempotent reader. Save the ID of each message in a table with a unique key. If it was already saved, do not do the work again.
Choosing the key and the number of Partitions
The key does two jobs at once: order and load spreading. These two sometimes fight each other.
- First ask where order is needed. In our shop, the order of messages of one order matters. The order between two different orders does not matter.
- Choose the smallest key that keeps the order you need. The order ID is better than the customer ID. A big customer with a thousand orders is spread over all Partitions when the key is the order.
- Leave some room for growth from the start. You can add Partitions later, but you cannot remove them.
- Adding Partitions changes the place of some keys. So at that moment the order of messages for those keys may get mixed up.
- Pick a number that divides well. For example, twelve Partitions split evenly between 2, 3, 4 or 6 pods.
Common mistakes
| Mistake | Result | Right way |
|---|---|---|
| Thinking order is kept across the whole topic | The “canceled” message is processed before “placed”. | Send related messages with one key. |
| More pods than Partitions | The extra pods stay idle and only cost money. | Keep the number of pods at most equal to the Partitions. |
| A key that is too coarse, like the customer ID | One Partition gets hot and the rest are idle. | A finer key, like the order ID. |
| Commit before processing | With every crash or deploy, a message is lost. | Process first, then Commit. |
| Thinking no duplicate message will come | The customer gets two emails or two text messages. | Build an Idempotent reader. |
| Saving to the database and then sending to Kafka, without Outbox | The order exists, but its message never arrives. | The Outbox pattern. |
| Ordering events by server clocks | Server clocks differ, and the order becomes wrong. | Partition order, or a version number from one source. |
When to use Kafka?
Good fit
- Several different services need the same events.
- The message volume is high.
- You must be able to read old messages again.
- The order of events for each entity (for example, each order) matters.
Poor fit
- You only need a simple work queue with a few workers. A queue like RabbitMQ is simpler.
- You need complex routing based on the message type.
- Each message must be confirmed or rejected on its own, without the others waiting behind it.
- The system is small, and the team cannot maintain a cluster.
Summary in seven lines
- Kafka is a log book. A message is not deleted after it is read.
- Each topic has several Partitions. Order is kept only inside one Partition.
- Same key means same Partition. Choose the key with care.
- Inside one Consumer Group, each Partition has only one reader.
- Readers beyond the number of Partitions stay idle.
- Process first, then Commit. So duplicate messages come, and the reader must be Idempotent.
- To send safely from the database to Kafka, use Outbox.