Part 4 · 1 chapters · ~8 min
Consumers and Consumer Groups
The poll loop, consumer groups and partition assignment, rebalancing (eager and cooperative, static membership), committing offsets (automatic and manual), at-least-once processing and idempotent consumers, consumer lag, scaling consumers, and handling poison messages with dead-letter topics.
6
Groups, offsets and lag
code
// manual commits after processing: at-least-once with an idempotent handler
await consumer.subscribe({ topics: ['transfers'] });
await consumer.run({
autoCommit: false,
eachMessage: async ({ topic, partition, message }) => {
const event = JSON.parse(message.value.toString());
await db.tx(async t => {
if (await t.oneOrNone('SELECT 1 FROM processed_events WHERE id = $1', [event.id])) return; // dedupe
await applyToReadModel(t, event);
await t.none('INSERT INTO processed_events (id) VALUES ($1)', [event.id]);
});
await consumer.commitOffsets([{ topic, partition, offset: (BigInt(message.offset) + 1n).toString() }]);
},
});
kafka-consumer-groups.sh --bootstrap-server b:9092 --describe --group notifier # LAG per partitionPoison messages (a record that always fails) block a partition forever if you keep retrying. Retry a few times, then publish it to a dead-letter topic with the error and move on; alert on DLT traffic.
CONSUMER GROUPS AND OFFSETS
partitions shared out among consumers; progress stored as committed offsets
swipe the figure sideways, or tap expand for full screen
1/5
sharing partitions
Consumers with the same group.id split the topic's partitions: each partition is read by exactly one consumer in the group. Two groups read the same data independently.
one consumer per partition per groupgroups are independent