Kembali ke jurnalCATATAN FAJAR
Software Engineering5 menit baca

At-Least-Once di Kafka: Retry dan Dead Letter Queue

Menguji timing commit di Kafka, retry yang dibatasi, pemrosesan duplikat, dan dead letter queue.

BAGIAN 6 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 Rebalancing
  5. 05Consumer Group Redis Streams: Competing Consumer dan PEL
  6. 06At-Least-Once di Kafka: Retry dan Dead Letter QueueKamu di sini
  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 6 bagian

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

SemanticCaraFailure mode
At-most-onceCommit offset sebelum memprosesCrash setelah commit, sebelum pekerjaan selesai: message hilang
At-least-onceCommit offset setelah memprosesCrash setelah pekerjaan selesai, sebelum commit: message diproses ulang (duplikat)
Exactly-onceOutput dan commit offset dalam satu transactionNggak 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.

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

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

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

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

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

ts
// 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.

TOPIK

SELANJUTNYA DI SERI INIRedis Streams: Message Nyangkut, XAUTOCLAIM, dan Dead Letter

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.