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:
XREADGROUP GROUP g c ... STREAMS key >mengirim entry baru ke consumercdan, di langkah yang sama, mencatat masing-masing entry di PEL sebagai milikc.- Consumer memproses entry itu dan memanggil
XACK, lalu entry-nya keluar dari PEL. - 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:
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:
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:
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:
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:
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.
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:
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.
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.
| Kafka | Redis Streams | |
|---|---|---|
| Pelacakan in-flight | satu committed offset per partition | satu entry PEL per message |
| Yang diblokir oleh crash | seluruh partition, sampai di-assign ulang | nggak ada yang lain; entry lain tetap jalan |
| Mekanisme retry | percobaan terbatas di dalam loop satu consumer | redelivery lewat beberapa sweep XAUTOCLAIM |
| Pemicu DLQ | penghitung percobaan lokal | delivery count dari XPENDING |
| Dead letter berakhir di | topic sampingan | stream 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.

Memuat komentar...