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 partition

Poison 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
partition 0partition 1partition 2partition 3consumer Ap0, p1consumer Bp2, p3consumer C joinsrebalance__consumer_offsetsgroup, partition → offset
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