Kembali ke jurnalCATATAN FAJAR
Software Engineering4 menit baca

Kafka Consumer Group: Lag, Draining, dan Rebalancing

Mengamati partition assignment, consumer lag, dan rebalancing saat member bergabung ke sebuah Kafka consumer group.

BAGIAN 4 DARI 16Kafka vs Redis Streams
  1. 01Kafka vs Redis Streams: Bikin Stock Ticker Dua Kali
  2. 02Kafka Partition dan Key: Menjaga Simbol Tetap Berurutan
  3. 03Dasar-Dasar Redis Streams: XADD, Entry ID, dan MAXLEN
  4. 04Kafka Consumer Group: Lag, Draining, dan RebalancingKamu di sini
  5. 05Consumer Group Redis Streams: Competing Consumer dan PEL
  6. 06At-Least-Once di Kafka: Retry dan Dead Letter Queue
  7. 07Redis Streams: Message Nyangkut, XAUTOCLAIM, dan Dead Letter
  8. 08Exactly-Once di Kafka: Idempotent Producer dan Transaction
  9. 09Redis Pub/Sub vs Streams: Ephemeral atau Durable
  10. 10Co-Partitioning: Key Sama, Partition Sama, di Kedua Topic
  11. 11Menulis Custom Partition Assigner KafkaJS untuk Local Join
  12. 12Kafka ISR vs Redis Cluster: Replication dan Failover
  13. 13Stateful Streams: Trigger Edge vs Level dan Window OHLC
  14. 14Membangun Live Market Dashboard dengan SSE dan Next.js
  15. 15Benchmark Kafka vs Redis Streams: Throughput dan Latency
  16. 16Kafka atau Redis Streams: Gimana Cara Memilihnya
Di artikel ini 5 bagian

Key di market-updates memberikan urutan per simbol. Consumer group membagi partition sebuah topic ke beberapa proses. Dengan begitu, lebih banyak consumer bisa berbagi beban dan memulihkan partition setelah sebuah member gagal.

Saya pengin paham perubahan kepemilikan selagi data terus berdatangan. Di eksperimen stock ticker saya, saya bikin backlog dan mengamati gimana backlog itu habis. Lalu saya menambahkan consumer ketiga untuk melihat partition mana yang pindah.

Yang dibagi itu partition, bukan message

Kafka consumer group membagi partition sebuah topic ke member-membernya. Setiap partition dimiliki oleh tepat satu consumer di group itu pada satu waktu. Dengan market-updates yang punya 6 partition:

  • 2 consumer di group: masing-masing 3 partition.
  • 3 consumer: masing-masing 2.
  • 6 consumer: masing-masing 1. Consumer ke-7 nggak kebagian partition apa pun dan cuma diam. Jumlah partition adalah batas atas yang mutlak untuk paralelisme yang berguna.

Aturan satu pemilik itulah yang bikin urutan per key tetap terjaga walaupun ada banyak consumer. Sebuah key, misalnya BBCA, selalu di-hash ke partition yang sama, dan partition itu punya tepat satu pemilik di setiap saat. Setiap tick BBCA ditangani oleh consumer yang sama, sesuai urutan tick itu diproduksi. Kepemilikan bisa pindah antar consumer, tapi urutan internal sebuah partition nggak pernah diatur ulang.

Test harness saya bikin ini bisa diamati dengan mendengarkan event GROUP_JOIN dari KafkaJS, yang terpicu di setiap member setiap kali group sepakat pada sebuah assignment:

typescript
function makeConsumer(label: string, assignments: Assignments): Consumer {
  const consumer = kafka.consumer({
    groupId: GROUP,
    sessionTimeout: 10000,
    heartbeatInterval: 3000,
  });
  consumer.on(consumer.events.GROUP_JOIN, (e) => {
    assignments.set(label, e.payload.memberAssignment);
  });
  return consumer;
}

heartbeatInterval mengatur frekuensi heartbeat. Coordinator pakai sessionTimeout untuk mendeteksi member yang udah nggak mengirim heartbeat. Handler yang block terlalu lama bisa bikin consumer kehilangan assignment-nya, walaupun prosesnya masih hidup.

Assignment yang dicetak memberikan tiga dari enam partition topic itu ke pod-1 dan tiga ke pod-2. Output-nya satu baris per member.

Committed offset dan apa yang diukur lag

Setiap partition menyimpan record dengan offset yang terus naik. Log end offset menunjukkan offset berikutnya setelah ujung log saat ini. Consumer group meng-commit offset tempat group itu harus melanjutkan. Setelah pemrosesan berhasil, biasanya ini adalah offset record berikutnya, bukan offset record terakhir yang diproses.

Selisih antara dua angka itu adalah lag:

text
lag = logEndOffset - committedOffset   (per partition, summed across the topic)

Lag nol berarti committed offset udah mencapai ujung partition yang dipakai dalam perhitungan. Lag positif yang stabil berarti masih ada backlog. Lag yang naik berarti backlog-nya bertambah. Helper saya menghitung lag dari metadata broker:

typescript
export async function lagWithAdmin(
  admin: Admin,
  groupId: string,
  topic: string,
): Promise<LagReport> {
  const [topicOffsets, groupOffsets] = await Promise.all([
    admin.fetchTopicOffsets(topic),
    admin.fetchOffsets({ groupId, topics: [topic] }),
  ]);

  const committed = new Map<number, number>();
  const entry = groupOffsets.find((g) => g.topic === topic);
  for (const p of entry?.partitions ?? []) committed.set(p.partition, Number(p.offset));

  const perPartition = topicOffsets.map((t) => {
    const high = Number(t.high);
    const low = Number(t.low);
    let c = committed.get(t.partition) ?? -1;
    if (c < 0) c = low; // group never committed here yet
    return { partition: t.partition, high, committed: c, lag: Math.max(0, high - c) };
  });

  return { total: perPartition.reduce((s, p) => s + p.lag, 0), perPartition };
}

Fallback c < 0 -> low menangani partition yang belum punya committed offset dari group. Fallback ini pakai offset paling awal yang masih disimpan, low, bukan nol. Ini mengukur history yang tersedia dengan asumsi group bakal mulai dari offset itu. Helper consumerLag() membuka dan menutup koneksi admin-nya sendiri.

Mengamati lag terkuras

Untuk bikin backlog yang kelihatan, saya memproduksi 200.000 tick ke 6 partition sebelum ada consumer yang terhubung. Producer-nya mengirim dalam chunk berisi 10.000 tick:

typescript
const CHUNK = 10000;
for (let sent = 0; sent < TOTAL; sent += CHUNK) {
  const msgs = sim.batch(Math.min(CHUNK, TOTAL - sent)).map(tickToKafka);
  await producer.send({ topic: TOPIC, messages: msgs });
}

Lalu pod-1 dan pod-2 bergabung ke group. Auto commit interval yang pendek bikin progres kelihatan lewat committed offset. Delay kecil di eachMessage bikin saya bisa melihat backlog-nya berkurang:

typescript
await c.run({
  autoCommitInterval: 400, // commit often so lag tracks real progress
  eachMessage: async ({ message }) => {
    const tick = tickFromKafka(message.value);
    for (const rule of SAMPLE_ALERTS) evaluate(rule, tick);
    counter.n++;
    if (counter.n % 1000 === 0) await sleep(20); // gentle throttle
  },
});

Selama itu jalan, saya polling consumerLag() setiap 400ms dan mencetak satu baris per sampel:

typescript
for (let i = 0; i < 25; i++) {
  const { total } = await consumerLag(GROUP, TOPIC);
  console.log(`t+${(i * 0.4).toFixed(1)}s   lag=${total}   consumed=${counter.n}`);
  if (total === 0 && counter.n >= TOTAL) break;
  await sleep(400);
}

Backlog-nya dimulai dengan 200.000 tick yang belum dikonsumsi sama sekali. Setiap consumer memproses 3 partition yang di-assign ke dirinya, dan lag turun seiring committed offset bergerak maju. Commit terakhir bisa terjadi setelah counter.n mencapai TOTAL. Kurva lag yang turun menunjukkan progres, sedangkan lag yang datar atau naik perlu diselidiki.

Rebalancing: apa yang terjadi waktu keanggotaan berubah

Setelah lag terkuras sampai nol, saya menambahkan consumer ketiga, pod-3, ke group yang sama dan membiarkannya stabil sebelum mencetak ulang assignment-nya:

typescript
const c3 = makeConsumer("pod-3", assignments);
await c3.connect();
await c3.subscribe({ topic: TOPIC, fromBeginning: true });
await c3.run({ eachMessage: async () => { counter.n++; } });
await sleep(3500);
printAssignments(assignments);

Panggilan subscribe dan run menambahkan member dan memicu rebalance. Assignment group berubah dari 3 partition per consumer jadi 2 partition per consumer. Kalau sebuah member terputus atau timeout, partition-nya pindah ke member yang tersisa.

mermaid
sequenceDiagram
    participant Pod1 as pod-1
    participant Pod2 as pod-2
    participant Pod3 as pod-3
    participant GC as "Group Coordinator"
    Note over Pod1,Pod2: "Owns 3 partitions each"
    Pod3->>GC: "Join the group"
    GC->>Pod1: "Revoke partitions, pause fetching"
    GC->>Pod2: "Revoke partitions, pause fetching"
    GC->>GC: "Recompute assignment across 3 members"
    GC->>Pod1: "New assignment: 2 partitions"
    GC->>Pod2: "New assignment: 2 partitions"
    GC->>Pod3: "New assignment: 2 partitions"
    Note over Pod1,Pod3: "GROUP_JOIN fires per member, fetching resumes"

Diagram itu menunjukkan rebalance "eager". Member yang udah ada berhenti mengonsumsi sementara selama group mengganti assignment-nya. Protokol cooperative atau incremental bisa membatasi partition yang pindah, tergantung dukungan broker dan client. Rebalance yang sering bakal mengganggu pekerjaan yang berguna, jadi pantau perubahan keanggotaan dan delay di handler.

Jaminan urutan tetap bertahan di tengah perubahan

Rebalance mengubah kepemilikan partition, tapi tetap menjaga urutan record yang tersimpan. Selama partitioning-nya nggak berubah, tick BBCA tetap masuk ke partition yang sama. Consumer baru melanjutkan dari committed offset, jadi consumer itu mungkin mengulang record yang udah diproses consumer sebelumnya tapi belum di-commit.

Consumer group di Redis Streams kelihatan mirip kalau dilihat sekilas dari luar. XREADGROUP juga menyebarkan pekerjaan ke sebuah group yang punya nama, tapi di baliknya nggak ada konsep partition dan sama sekali nggak ada kepemilikan per key. Consumer mana pun bisa mengambil entry mana pun, termasuk dua entry untuk simbol yang sama di waktu yang sama.

Redis menangani delivery yang belum selesai lewat pending entries list. Post berikutnya membandingkan model recovery itu dengan partition reassignment di Kafka.

TOPIK

SELANJUTNYA DI SERI INIConsumer Group Redis Streams: Competing Consumer dan PEL

MAKASIH UDAH BACA

Gimana menurutmu?

Reaksi atau obrolan, dua-duanya selalu ditunggu.

Memuat reaksi…

Bagikan

Memuat komentar...

LANJUT JELAJAH

Satu pikiran bawa ke pikiran lain.

Semua tulisan
Kembali ke semua tulisanSatu catatan, pelan-pelan.