"At-least-once delivery" bergantung pada commit offset setelah handling berhasil. Consumer tetap bisa mengulang pekerjaan kalau crash setelah memproses tapi sebelum commit. Handler-nya harus membuat efek yang berulang tetap aman.
Redis menyimpan entry yang udah dikirim di Pending Entries List sampai di-acknowledge. Kafka mencatat committed offset per group dan partition. Aplikasi yang memutuskan kapan handling yang berhasil mengizinkan offset itu maju.
Saya menguji consumer yang crash, rebalance di tengah batch, dan message yang gagal di setiap percobaan. Post ini menjelaskan perilaku retry dan dead letter di Kafka.
Delivery semantics adalah keputusan commit
Kafka nggak melacak "did the consumer finish processing this." Kafka melacak satu angka per partition: committed offset. Semua hal soal delivery semantics ditentukan oleh kapan kamu menggeser angka itu.
| Semantic | Cara | Failure mode |
|---|---|---|
| At-most-once | Commit offset sebelum memproses | Crash setelah commit, sebelum pekerjaan selesai: message hilang |
| At-least-once | Commit offset setelah memproses | Crash setelah pekerjaan selesai, sebelum commit: message diproses ulang (duplikat) |
| Exactly-once | Output dan commit offset dalam satu transaction | Nggak ada yang hilang maupun duplikat, nggak ada efek parsial |
Commit sebelum memproses bisa melewatkan pekerjaan yang belum selesai setelah crash. Commit setelah memproses memungkinkan replay dari offset sebelumnya. Consumer berikutnya lalu bisa mengulang sebuah efek, misalnya notifikasi alert.
Kafka transaction bisa meng-commit output Kafka dan source offset bersamaan, seperti yang dijelaskan di post tentang transaction. Tabel di atas mengasumsikan input yang masih di-retain dan consumer yang dikonfigurasi dengan benar. Jaminan transaction-nya mencakup record dan offset Kafka, bukan efek eksternal.
Commit sesudah, bukan sebelum
Di KafkaJS, matikan auto commit dan commit secara manual setelah message selesai di-handle.
await consumer.run({
autoCommit: false,
eachMessage: async ({ topic, partition, message }) => {
await handle(message);
await consumer.commitOffsets([
{ topic, partition, offset: (Number(message.offset) + 1).toString() },
]);
},
});
Posisi commitOffsets relatif terhadap handle menentukan titik recovery. Kalau commit terjadi duluan, recovery normal dimulai setelah record itu walaupun handling-nya gagal. Commit setelah handling berhasil memungkinkan replay kalau consumer gagal sebelum commit. Aplikasinya harus tahan terhadap handling yang berulang.
Retry tanpa memblokir partition
Commit offset setelah memproses mendukung recovery dari crash, tapi sebuah message bisa gagal di setiap percobaan. Kalau consumer nggak pernah maju melewatinya, record berikutnya di partition itu nggak bisa diproses. Satu record yang rusak bisa menunda setiap symbol di partition itu.
Saya membuat kegagalannya deterministik supaya bisa menguji perilaku retry. Consumer mengklasifikasikan tick di market-updates berdasarkan sequence number untuk tiap symbol. Sebagian besar langsung berhasil. Sebagian gagal di beberapa percobaan pertama, lalu berhasil. Poison tick gagal di setiap percobaan:
function classify(t: Tick): Kind {
if (t.seq % 25 === 0) return "poison"; // permanent failure -> belongs in the DLQ
if (t.seq % 7 === 0) return "flaky"; // transient -> succeeds once retried
return "ok";
}
Di dalam eachMessage, sebuah message dapat jatah percobaan lokal yang terbatas, dengan backoff singkat di antaranya, sebelum consumer menyerah:
let ok = false;
let lastErr: unknown;
for (let attempt = 1; attempt <= MAX_ATTEMPTS; attempt++) {
try {
processTick(tick, attempt);
ok = true;
break;
} catch (e) {
lastErr = e;
await sleep(15); // backoff before retry
}
}
Ini adalah retry lokal: partition tetap berada di record yang sama selama setiap percobaan. Kegagalan downstream yang sementara mungkin pulih sebelum batas retry habis. Batas MAX_ATTEMPTS mencegah message yang selalu gagal memblokir partition tanpa batas waktu.
Dead letter queue
Setelah batas retry tercapai, consumer menulis record yang gagal ke topic dead letter queue (DLQ). Record itu menyertakan konteks yang cukup untuk investigasi nanti. Setelah write ke DLQ berhasil, consumer bisa memajukan source offset-nya.
graph TD
M["Message arrives"] --> P["processTick(tick, attempt)"]
P -->|success| C["commitOffsets"]
P -->|failure, attempt < MAX| R["backoff, retry"]
R --> P
P -->|failure, attempt = MAX| D["dlqProducer.send to DLQ topic"]
D --> C
Write ke DLQ membawa key dan value aslinya, ditambah header yang menjelaskan kenapa gagal dan dari mana asalnya:
await dlqProducer.send({
topic: DLQ,
messages: [
{
key: message.key,
value: message.value,
headers: {
error: lastErr instanceof Error ? lastErr.message : String(lastErr),
origin: `${topic}/${partition}/${message.offset}`,
},
},
],
});
Yang terjadi tepat setelahnya itu penting: consumer meng-commit offset aslinya, baik message itu berhasil maupun masuk dead letter.
// Commit AFTER handling => at-least-once.
await consumer.commitOffsets([
{ topic, partition, offset: (Number(message.offset) + 1).toString() },
]);
Source offset harus maju setelah write ke DLQ berhasil, kalau nggak consumer bakal membaca record gagal yang sama lagi. Saya memakai createIdempotentProducer() untuk producer DLQ. Idempotence di producer menangani retry dari request-nya sendiri. Ini nggak mencegah munculnya record DLQ lain kalau terjadi crash di antara write ke DLQ dan commit source offset.
Saya memberi topic DLQ saya satu partition aja. Eksperimen ini nggak butuh pemrosesan dead letter secara paralel. Topic itu menyimpan tick yang gagal untuk diperiksa. Setelah ada fix, operator bisa memutuskan mau me-replay atau membuangnya.
Idempotent consumer: jawaban yang praktis
At least once delivery tetap bisa mengulang pemrosesan walaupun udah ada retry yang dibatasi dan DLQ. Rebalance atau crash setelah handle() tapi sebelum commitOffsets bisa menyebabkan duplikat. Retry lokal juga bisa mengulang sebagian pekerjaan. Consumer harus membuat efeknya idempotent.
Pakai identifier message yang stabil, misalnya symbol ditambah sequence number. Untuk efek ke database, terapkan uniqueness pada (symbol, seq) dan pakai ON CONFLICT DO NOTHING kalau cocok dengan operasinya. Simpan keputusan deduplikasi dan efek database dalam transaction yang sama.
Notifikasi eksternal butuh dukungan idempotency atau protokol pengiriman sendiri. Pengecekan terpisah yang diikuti pengiriman bisa kena race atau gagal di antara kedua operasi itu. Pemrosesan yang berulang harus tetap menghasilkan hasil yang dimaksud.
Saya memakai kata yang sama untuk dua jenis idempotence di kode saya. createIdempotentProducer() mencegah record duplikat akibat retry di producer. Kafka memakai producer ID dan sequence number untuk tiap partition untuk mengenali pengiriman yang berulang.
Commit offset yang telat tetap bisa membuat consumer memproses record yang sama lagi. Logic consumer harus menangani sumber duplikat yang terpisah ini.
Yang bakal saya ship
Saya memakai delivery yang mungkin di-retry dan idempotent consumer sebagai default. Crash dan rebalance bisa mengulang pekerjaan, jadi efek yang berulang harus aman. Kafka transaction cocok untuk pipeline yang output dan offset-nya sama-sama tetap di Kafka. Kafka transaction nggak bisa membuat API pembayaran atau order eksternal jalan tepat sekali.
Pantau DLQ dan kirim alert waktu ada record yang masuk. Operator butuh cara untuk memeriksa kegagalan dan memutuskan apakah tiap record di-replay atau dibuang. Proses recovery ini menjaga pekerjaan yang gagal tetap terlihat setelah partition utama maju.
Redis Streams menyediakan PEL, delivery count, dan XAUTOCLAIM. Aplikasinya harus menggabungkan semua itu dengan batas retry dan policy dead letter. Post berikutnya menjelaskan implementasinya.

Memuat komentar...