Kembali ke jurnalCATATAN FAJAR
Software Engineering5 menit baca

Consumer Group Redis Streams: Competing Consumer dan PEL

Pakai XREADGROUP, Pending Entries List, dan XACK untuk membagikan entry stream Redis dan melacak pekerjaan yang belum selesai.

BAGIAN 5 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 PELKamu di sini
  6. 06At-Least-Once di Kafka: Retry dan Dead Letter Queue
  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 8 bagian

Kafka dan Redis Streams sama-sama menyediakan "consumer group," tapi unit kerja yang mereka bagikan berbeda. Kafka membagikan partition ke member group. Redis Streams membagikan entry yang tersedia ke consumer yang memintanya. Perbedaan ini berpengaruh ke urutan dan recovery.

Consumer group Kafka membagikan partition ke member dan mengubah assignment ketika keanggotaannya berubah. Redis memakai XREADGROUP, Pending Entries List, dan XACK untuk membagikan dan melacak entry. Saya membandingkan keduanya sambil membangun stock ticker yang sama dua kali. Kedua versi itu butuh perilaku yang mirip di bawah beban, walaupun model delivery-nya berbeda.

Dari partition ke work queue

Di contoh Kafka ini, dua consumer masing-masing menerima tiga dari enam partition. Key routing yang stabil menaruh tick BBCA di partition yang sama. Consumer yang di-assign ke partition itu menangani tick-tick tersebut sampai sebuah rebalance mengubah kepemilikannya.

XREADGROUP di Redis dengan ID > mengembalikan entry yang belum pernah dikirim oleh group itu. Member mana pun bisa meminta batch berikutnya. Karena itu, dua entry untuk simbol yang sama bisa sampai ke consumer yang berbeda dan selesai dengan urutan yang nggak sesuai.

Consumer yang lebih banyak bisa berbagi pekerjaan tanpa batas jumlah partition, tapi kapasitas Redis dan workload-nya tetap membatasi concurrency yang berguna. Urutan untuk satu simbol butuh desain tersendiri, misalnya satu stream dan satu worker aktif per simbol.

Membuat group: XGROUP CREATE

Saya membuat group untuk stream key yang menerima tick lewat XADD:

text
XGROUP CREATE market-updates evaluators 0 MKSTREAM

Argumen sebelum MKSTREAM menentukan posisi awal group. $ mulai setelah entry yang udah ada waktu group dibuat. 0 mulai sebelum entry pertama yang masih disimpan. MKSTREAM membuat stream-nya kalau belum ada. Tanpa opsi itu, pembuatan group butuh stream yang udah ada.

Saya membungkusnya dalam sebuah helper yang menelan error "group already exists", jadi kode setup saya bisa memanggilnya di setiap run tanpa perlu dipikirin:

typescript
export async function ensureGroup(
  redis: Redis,
  key: string,
  group: string,
  start = "$",
): Promise<void> {
  try {
    await redis.xgroup("CREATE", key, group, start, "MKSTREAM");
  } catch (e) {
    if (!(e instanceof Error) || !e.message.includes("BUSYGROUP")) throw e;
  }
}

Saya memberikan start = "0" sebelum menulis entry apa pun. Sebanyak 200 tick yang dibuat setelahnya membentuk backlog group yang belum dibaca.

XREADGROUP: consumer yang pertama minta, dia yang dapat

Pembacaan group butuh nama group dan nama consumer. Berbeda dengan XREAD biasa, pembacaan ini juga mengubah state delivery milik group:

text
XREADGROUP GROUP evaluators pod-1 COUNT 10 STREAMS market-updates >

Nama consumer pod-1 menandai sebuah member group. Redis nggak memverifikasi bahwa nama itu cocok dengan satu proses yang unik, jadi pakai nama yang berbeda untuk consumer yang independen. Setiap panggilan XREADGROUP > meminta batch berikutnya yang tersedia. Nggak ada rebalance partition sebelum delivery.

Dibaca bukan berarti di-acknowledge: pending entries list

Untuk delivery XREADGROUP > yang biasa tanpa NOACK, Redis mencatat entry-nya di Pending Entries List (PEL). Metadata pending-nya mencakup entry ID, nama consumer, waktu delivery, dan jumlah delivery. Pemrosesan yang berhasil di aplikasi aja nggak menghapus state pending ini.

Dibaca dan di-acknowledge adalah dua state yang berbeda. Mendapatkan entry dari XREADGROUP cuma berarti Redis menandainya sebagai udah dikirim ke kamu. Itu sama sekali nggak bilang apakah kode kamu udah melakukan sesuatu dengan entry itu.

Kamu bisa memeriksa PEL secara langsung:

text
XPENDING market-updates evaluators

Ringkasannya mencakup jumlah pending, ID pending terendah dan tertinggi, dan jumlah untuk setiap consumer. Bentuk extended-nya menampilkan idle time dan jumlah delivery setiap entry. Nilai-nilai ini membantu proses recovery memilih entry yang mau di-reclaim. Post tentang recovery menjelaskan XAUTOCLAIM.

XACK: memberi tahu Redis kalau kamu udah selesai

Setelah penanganannya berhasil, pakai XACK untuk menghapus record pending-nya. Acknowledgment ini nggak menghapus entry di stream:

text
XACK market-updates evaluators 1720000000000-0 1720000000001-0

Drain loop saya menaruh seluruh siklus baca lalu ack di satu tempat:

typescript
async function drain(
  conn: Redis,
  consumer: string,
  counters: Record<string, number>,
): Promise<void> {
  for (;;) {
    const entries = await readGroup(conn, { key: KEY, group: GROUP, consumer, count: 10 });
    if (entries.length === 0) break;
    counters[consumer] = (counters[consumer] ?? 0) + entries.length;
    await conn.xack(KEY, GROUP, ...entries.map((e) => e.id));
  }
}

Setelah readGroup, consumer menghitung entry-nya dan meng-acknowledge setiap ID lewat xack. Crash di antara kedua panggilan ini bakal meninggalkan entry-nya di PEL. Consumer lain bisa me-reclaim entry itu, atau consumer aslinya bisa membaca entry pending miliknya setelah pulih.

mermaid
graph TD
    A["XADD market-updates"] --> B[("Stream entries")]
    B -->|"XREADGROUP GROUP evaluators pod-1 >"| C["pod-1 processes"]
    B -->|"XREADGROUP GROUP evaluators pod-2 >"| D["pod-2 processes"]
    C --> E[("PEL entry, owner=pod-1")]
    D --> F[("PEL entry, owner=pod-2")]
    C -->|"XACK"| G["Removed from PEL"]
    D -->|"XACK"| H["Removed from PEL"]

Kedua consumer menarik data dari node stream yang sama, bukan dari partition terpisah yang di-assign ke masing-masing, dan di situlah model ini berbeda dari Kafka.

Dua consumer, tanpa kepemilikan

Tesnya memakai group bernama evaluators dan 200 tick hasil simulasi. Dua koneksi ioredis menamai consumer-nya pod-1 dan pod-2. Keduanya menjalankan drain loop lewat Promise.all:

mermaid
sequenceDiagram
    participant P1 as pod-1
    participant R as Redis
    participant P2 as pod-2
    P1->>R: "XREADGROUP GROUP evaluators pod-1 COUNT 10 STREAMS market-updates >"
    R-->>P1: "10 entries, PEL += 10, owner=pod-1"
    P2->>R: "XREADGROUP GROUP evaluators pod-2 COUNT 10 STREAMS market-updates >"
    R-->>P2: "next 10 entries, PEL += 10, owner=pod-2"
    Note over R: "PEL now holds 20 unacked entries"
    P1->>R: "XACK market-updates evaluators ...ids"
    P2->>R: "XACK market-updates evaluators ...ids"
    Note over R: "PEL empty, all 20 acknowledged"

Setiap panggilan readGroup meminta sampai 10 entry. Kedua consumer berebut 200 entry itu, jadi pembagiannya bergantung pada timing request. Setelah kedua loop selesai, pendingSummary melaporkan nol entry pending karena kedua loop udah meng-acknowledge setiap entry yang dihitung.

Consumer group Redis nggak membagikan simbol ke consumer. Consumer yang berbeda bisa memproses tick BBCA yang berurutan secara bersamaan dan selesai dengan urutan yang berbeda. Tes saya cuma menghitung entry, jadi hal ini nggak mengubah hasilnya.

Ticker-nya mengevaluasi perubahan harga, dan di situ urutan memengaruhi hasil. Field simbol di market-updates nggak memaksakan kepemilikan consumer. Stream terpisah dengan assignment consumer yang terkontrol, atau partition Kafka, bisa memberikan urutan pemrosesan yang dibutuhkan.

Kapan work queue lebih unggul dari kepemilikan partition

Group di Redis memungkinkan worker tambahan meminta entry tanpa reassignment partition. Ini cocok untuk job independen yang bisa selesai dalam urutan apa pun. Throughput tetap bergantung pada kapasitas Redis, biaya handler, dan resource bersama lainnya. Consumer yang lebih banyak nggak menjamin throughput yang lebih tinggi.

Group KafkaGroup Redis
AssignmentPartition dipatok ke memberEntry diberikan ke siapa pun yang minta
Urutan per key di antara consumerTerjagaNggak terjaga
Paralelisme maksimum yang bergunaSama dengan jumlah partitionNggak terbatas, sampai Redis jadi bottleneck
Menambah kapasitasTambah partition dan consumer, memicu rebalanceTinggal jalankan consumer tambahan
Dibaca vs. selesaiOffset yang di-commitEntry PEL dihapus oleh XACK
Sinyal backpressureConsumer lagUkuran PEL dan panjang stream

Tradeoff assignment

Group Kafka membagi partition ke member-membernya, dengan urutan di dalam setiap partition. Group Redis membagikan entry satu per satu, jadi worker yang jalan bersamaan bisa menyelesaikan entry yang saling berkaitan dengan urutan yang nggak sesuai. Redis mencatat entry yang udah dikirim di PEL sampai acknowledgment atau operasi lain menghapus state pending-nya.

Entry yang butuh acknowledgement tetap pending sampai group menghapus record pending tersebut. Kode recovery harus memutuskan kapan entry itu di-retry atau dikirim ke dead letter stream. Post berikutnya membandingkan timing commit, retry yang dibatasi, dan dead letter di Kafka.

TOPIK

SELANJUTNYA DI SERI INIAt-Least-Once di Kafka: Retry dan Dead Letter Queue

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.