Kembali ke jurnalCATATAN FAJAR
Software Engineering6 menit baca

Redis Streams: Message Nyangkut, XAUTOCLAIM, dan Dead Letter

Memeriksa entry Redis yang pending, mengambil alih delivery yang idle dengan XAUTOCLAIM, dan menangani kegagalan berulang dengan dead letter stream.

BAGIAN 7 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 Queue
  7. 07Redis Streams: Message Nyangkut, XAUTOCLAIM, dan Dead LetterKamu di sini
  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 7 bagian

Setelah sebuah Kafka consumer gagal, member lain bisa melanjutkan partition-nya dari committed offset. Recovery di Redis Streams bekerja per entry. Consumer yang gagal setelah XREADGROUP bisa meninggalkan pekerjaan di pending entries list. Aplikasinya harus memakai sebuah proses recovery, misalnya XAUTOCLAIM, untuk mengambil pekerjaan yang belum selesai itu.

Post sebelumnya menguji retry Kafka yang dibatasi. Sebuah poison tick di market-updates pada akhirnya pindah ke dead letter topic. Saya menguji kegagalan yang sama di Redis. Proses recovery-nya memakai pending entry dan delivery count, bukan partition offset.

Pending entries list

Consumer group di Redis Streams melacak pekerjaan yang masih in flight dengan sebuah struktur bernama PEL (pending entries list), satu per stream dan group. Mekanismenya sederhana:

  1. XREADGROUP GROUP g c ... STREAMS key > mengirim entry baru ke consumer c dan, di langkah yang sama, mencatat masing-masing entry di PEL sebagai milik c.
  2. Consumer memproses entry itu dan memanggil XACK, lalu entry-nya keluar dari PEL.
  3. Atau consumer-nya nggak pernah memanggil XACK, karena crash, hang, atau mati gara-gara deploy. Entry itu tetap pending atas nama consumer tersebut.

Awalnya saya mengira ada visibility timeout kayak di SQS. Redis nggak otomatis membuat entry yang pending tersedia untuk pembaca baru setelah timeout. Sebuah proses recovery harus meminta entry yang pending atau mengambil alihnya.

Saya membuat group evaluators di sebuah stream berisi market tick. Consumer pod-crash membaca entry tapi nggak meng-acknowledge satu pun:

ts
await ensureGroup(redis, KEY, GROUP, "0");

const sim = new MarketSimulator({ seed: 33, volatilityScale: 3 });
for (const tick of sim.batch(TOTAL)) await redis.xadd(KEY, "*", ...tickToRedis(tick));

// "pod-crash" reads everything but ACKs nothing (simulating a crash).
const taken = await readGroup(redis, {
  key: KEY,
  group: GROUP,
  consumer: "pod-crash",
  count: TOTAL,
});

ensureGroup adalah wrapper tipis di atas XGROUP CREATE ... MKSTREAM, dan readGroup menyusun panggilan XREADGROUP yang variadic. Dengan TOTAL sebanyak 50 tick untuk 10 simbol, ke-50 entry itu sekarang ada di PEL, semuanya milik pod-crash, dan pod-crash nggak pernah kembali.

XPENDING: siapa memegang apa

Sebelum bisa mengambil alih apa pun, kamu harus bisa melihatnya dulu. XPENDING punya dua bentuk, dan akhirnya saya bikin wrapper untuk keduanya.

Bentuk ringkasan (tanpa id range) memberi kamu jumlah total yang pending dan rinciannya per consumer:

ts
export async function pendingSummary(redis: Redis, key: string, group: string) {
  const res = await redis.call("XPENDING", key, group) as [
    number,
    string | null,
    string | null,
    Array<[string, string]> | null,
  ];
  const perConsumer: Record<string, number> = {};
  for (const [c, n] of res?.[3] ?? []) perConsumer[c] = Number(n);
  return { total: Number(res?.[0] ?? 0), perConsumer };
}

Tepat setelah pod-crash membaca lalu menelantarkan 50 entry-nya, ini mencetak now PENDING: 50 (owned by pod-crash). Bentuk itulah sinyal crash yang bakal kamu temui di dunia nyata: nama consumer yang muncul di perConsumer dan nggak pernah menyusut.

Bentuk extended menerima rentang ID dan sebuah count. Setiap hasil berisi ID entry, pemiliknya, idle time, dan delivery count:

ts
export async function pendingDetails(redis: Redis, key: string, group: string, count = 100) {
  const res = await redis.call("XPENDING", key, group, "-", "+", String(count)) as Array<
    [string, string, number, number]
  >;
  return (res ?? []).map(([id, consumer, idleMs, deliveries]) => ({
    id,
    consumer,
    idleMs: Number(idleMs),
    deliveries: Number(deliveries),
  }));
}

Field deliveries itulah yang penting. Field ini adalah jawaban bawaan Redis untuk "how many times has this specific entry been handed to a consumer," dan keputusan dead letter seharusnya didasarkan pada field ini.

XAUTOCLAIM: mengambil alih pekerjaan yang idle

XAUTOCLAIM adalah alat yang memindahkan kepemilikan. Command ini menerima minimum idle time dan memilih entry pending yang idle time-nya memenuhi batas itu:

ts
export async function autoClaim(
  redis: Redis,
  p: { key: string; group: string; consumer: string; minIdleMs: number; start?: string; count?: number },
) {
  const res = await redis.call(
    "XAUTOCLAIM",
    p.key,
    p.group,
    p.consumer,
    String(p.minIdleMs),
    p.start ?? "0-0",
    "COUNT",
    String(p.count ?? 100),
  ) as [string, Array<[string, string[]]>, string[]?];
  return { cursor: res[0], claimed: parseEntries(res[1]) };
}

Minimum idle time memilih entry yang udah cukup lama pending untuk diselidiki. Nilai ini nggak bisa membuktikan bahwa consumer-nya gagal. Consumer yang lambat mungkin masih memproses sebuah entry setelah consumer lain meng-claim entry itu. Tentukan threshold-nya dari perkiraan waktu pemrosesan, dan pastikan efek yang berulang tetap aman.

Saya membiarkan entry-entry itu diam sebentar, lalu mengambil alihnya dengan threshold yang lebih pendek dari waktu tunggunya:

ts
await sleep(150); // let entries accrue idle time

const { claimed } = await autoClaim(redis, {
  key: KEY,
  group: GROUP,
  consumer: "pod-rescue",
  minIdleMs: 100,
  count: TOTAL,
});

Setelah 150ms, entry-entry itu melewati idle threshold 100ms. Tesnya mengambil ke-50 entry itu atas nama pod-rescue. Kepemilikan berpindah secara atomic di dalam Redis, tapi consumer sebelumnya nggak di-fence dari pekerjaan eksternal. Jadi pemrosesannya harus tahan terhadap percobaan yang tumpang tindih.

mermaid
sequenceDiagram
    participant PC as pod-crash
    participant R as Redis
    participant PR as pod-rescue
    PC->>R: XREADGROUP GROUP evaluators pod-crash COUNT 50
    Note over R: 50 entries enter the PEL, owned by pod-crash
    Note over PC: pod-crash dies before XACK
    Note over R: entries sit idle, still owned by pod-crash
    PR->>R: XAUTOCLAIM key evaluators pod-rescue 100 0-0 COUNT 50
    Note over R: idle >= 100ms, so ownership moves to pod-rescue
    R-->>PR: 50 claimed entries, delivery count incremented

Catatan terakhir itu penting: meng-claim sebuah entry lewat XAUTOCLAIM dihitung sebagai satu delivery, sama seperti XREADGROUP yang pertama. Setelah satu kali reclaim ini, setiap entry dari ke-50 entry itu udah menunjukkan delivery count 2.

Dari delivery count ke dead letter

Saya memakai delivery count dari XPENDING untuk memilih entry yang masuk DLQ. Setelah N kali delivery, handler menyalin entry itu ke sebuah DLQ stream dan memanggil XACK untuk entry aslinya.

Contoh Kafka menghitung percobaan di dalam satu loop consumer. Versi Redis ini menghitung delivery yang tercatat lintas consumer. Sebuah delivery nggak membuktikan bahwa handler-nya menyelesaikan satu percobaan.

Sebuah recovery loop meng-claim entry yang idle, mengecek deliveries lewat pendingDetails, dan memproses entry yang masih di bawah batas retry. Entry yang udah mencapai batas masuk ke DLQ. Acknowledgment dilakukan setelah pemrosesan berhasil atau setelah write ke DLQ berhasil.

Tes ini menjalankan satu siklus recovery, jadi setiap entry yang di-claim punya delivery count 2. Angka itu nggak bisa membedakan kegagalan yang terus-menerus dari entry yang ditelantarkan pod-crash. Karena itu, untuk demonstrasi ini saya memakai seq sebagai penanda poison yang deterministik:

ts
let acked = 0;
let deadLettered = 0;
for (const entry of claimed) {
  const isPoison = Number(entry.fields.seq) === 3; // ~1 per symbol -> ~10 poison
  if (isPoison) {
    await redis.xadd(DLQ, "*", ...Object.entries(entry.fields).flat(), "deadLetteredFrom", entry.id);
    await redis.xack(KEY, GROUP, entry.id); // ack so it leaves the PEL
    deadLettered++;
  } else {
    await redis.xack(KEY, GROUP, entry.id);
    acked++;
  }
}

Write ke DLQ menyalin field aslinya dan menambahkan deadLetteredFrom berisi ID entry aslinya. Pemrosesan yang berhasil maupun write ke DLQ yang berhasil sama-sama diikuti XACK.

Write ke DLQ dan acknowledgment yang terpisah punya celah kalau terjadi crash. Restart bisa menulis entry DLQ lain sebelum acknowledgment berhasil. Pakai idempotency atau operasi Redis yang atomic kalau key yang dibutuhkan memungkinkan.

Run ini memakai 50 entry untuk 10 simbol, dengan kira-kira satu poison tick per simbol. Hasil akhirnya sekitar 40 diproses, sekitar 10 dikirim ke DLQ, dan PENDING di angka 0. Pengecekan pendingDetails terakhir mencetak entry pending yang masih tersisa. Nggak ada yang dicetak kalau PEL-nya kosong.

mermaid
graph TD
    A["XAUTOCLAIM reclaims an idle entry"] --> B["XPENDING: check deliveries"]
    B -->|"deliveries <= threshold"| C["process, then XACK"]
    B -->|"deliveries > threshold"| D["XADD to DLQ stream"]
    D --> E["XACK original id"]
    C --> F["entry leaves the PEL"]
    E --> F

Satu poison entry nggak memblokir yang lain

Kafka group memberikan setiap partition ke satu member. Kalau handler-nya berhenti di sebuah poison record, record berikutnya di partition itu harus menunggu. Redis group melacak setiap entry pending secara terpisah. XREADGROUP > bisa terus mengirim entry baru selagi entry sebelumnya masih pending.

KafkaRedis Streams
Pelacakan in-flightsatu committed offset per partitionsatu entry PEL per message
Yang diblokir oleh crashseluruh partition, sampai di-assign ulangnggak ada yang lain; entry lain tetap jalan
Mekanisme retrypercobaan terbatas di dalam loop satu consumerredelivery lewat beberapa sweep XAUTOCLAIM
Pemicu DLQpenghitung percobaan lokaldelivery count dari XPENDING
Dead letter berakhir ditopic sampinganstream sampingan

Saya melihat tradeoff urutan ini waktu pertama kali menguji consumer group Redis. Dua consumer bisa memproses tick untuk simbol yang sama secara bersamaan. Tick yang belakangan bisa selesai duluan. Delivery entry yang independen menghindari blocking di seluruh stream, tapi aplikasinya harus menjaga sendiri urutan pemrosesan yang dibutuhkan.

Delivery berulang tetap mungkin terjadi

Redis Streams bisa mendukung delivery berulang kalau aplikasinya mengambil alih pekerjaan yang pending. Baik entry PEL maupun XACK nggak otomatis menjalankan recovery. Retention dan persistence juga harus menjaga entry-nya tetap ada.

Lua script atau MULTI/EXEC bisa menggabungkan command output dan acknowledgment. Redis Cluster mensyaratkan key yang terlibat berada di slot yang sama. Efek eksternal butuh koordinasi terpisah, dan error command Redis nggak memberikan rollback yang transactional.

Bikin handler-nya idempotent. pod-crash mungkin udah menyelesaikan sebagian pekerjaan sebelum consumer lain mengambil alih entry-nya. pod-rescue harus tahan terhadap pemrosesan berulang. Pakai identitas yang stabil, misalnya (symbol, seq) untuk sebuah tick atau ID notifikasi.

Yang bakal saya pantau di production

Pantau pendingSummary untuk entry yang terus pending. Jumlah PENDING yang naik bisa menandakan consumer yang lambat, gagal, atau nggak aktif. Itu nggak membuktikan ada crash. Periksa idle time dan kesehatan consumer sebelum recovery.

Set idle threshold XAUTOCLAIM di atas perkiraan waktu pemrosesan, tapi anggap aja consumer yang lambat masih bisa tumpang tindih dengan recovery. Pakai delivery count dan konteks kegagalannya untuk membatasi retry. Pantau DLQ-nya dan sediakan proses untuk inspeksi dan replay.

Kedua implementasi ticker butuh handler yang idempotent kalau pekerjaan bisa berulang. Post berikutnya membahas idempotence di producer Kafka dan cakupan Kafka transaction.

TOPIK

SELANJUTNYA DI SERI INIExactly-Once di Kafka: Idempotent Producer dan Transaction

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.