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:
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:
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:
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:
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:
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:
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:
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.
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.

Memuat komentar...