Kembali ke jurnalCATATAN FAJAR
Software Engineering5 menit baca

Redis Pub/Sub vs Streams: Ephemeral atau Durable

Membandingkan delivery Pub/Sub yang langsung dengan retention stream dan pembacaan belakangan, termasuk batasan persistence Redis.

BAGIAN 9 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 Letter
  8. 08Exactly-Once di Kafka: Idempotent Producer dan Transaction
  9. 09Redis Pub/Sub vs Streams: Ephemeral atau DurableKamu di sini
  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

Redis menyediakan Pub/Sub lewat PUBLISH dan entry stream yang disimpan lewat XADD. Saya menganggap keduanya bisa saling menggantikan, sampai sebuah subscriber ticker terputus selama sekitar satu detik. Pub/Sub nggak bisa memutar ulang pesan-pesan yang terlewat. Stream bisa mengembalikan entry yang masih disimpan.

Seri stock ticker membandingkan Kafka dan Redis Streams. Post ini membandingkan dua fitur Redis: Pub/Sub dan Streams.

Broker yang sama, dua janji yang nggak berhubungan

PUBLISH dan SUBSCRIBE adalah fitur messaging asli Redis, bertahun-tahun lebih tua dari Streams. Publisher memanggil PUBLISH channel message, dan Redis mendorong pesan itu ke setiap client yang saat itu subscribe ke channel. Nilai kembaliannya berupa angka: berapa banyak subscriber yang berhasil dijangkau.

XADD adalah command yang saya bedah di Redis Streams Fundamentals. Command ini menambahkan data ke sebuah log yang tinggal di dalam sebuah key Redis. Setiap entry dapat ID dengan format <millisecondsTime>-<sequence>, dan entry-nya tetap ada di stream sampai ada yang men-trim-nya, bukan sampai ada yang membacanya.

Yang membedakan keduanya adalah apakah pesannya tetap ada setelah dikirim. Replay, consumer group, at least once delivery: semuanya mengikuti dari jawaban itu.

Pub/Sub: kamu dapat apa yang lagi mendengarkan

Pub/Sub mengirim pesan ke client yang saat itu subscribe ke channel. Dia nggak menyimpan history untuk subscriber yang datang belakangan. Kalau nggak ada subscriber, PUBLISH mengembalikan nol dan pesannya nggak disimpan.

Tes paling kecil yang kepikiran sama saya: publish ke channel yang belum ada subscriber-nya, lalu lihat apa yang dikembalikan.

typescript
const CHANNEL = "demo:s04:ticks";

// Pub/Sub - publish with no subscriber
const receivers0 = await pub.publish(CHANNEL, "BBCA@9500");
console.log(`PUBLISH reached ${receivers0} subscribers -> message LOST`);
// receivers0 is 0. The message does not exist anywhere any more.

Panggilan yang sama sekali lagi, kali ini setelah sebuah subscriber terhubung:

typescript
const received: string[] = [];
sub.on("message", (_ch, msg) => received.push(msg));
await sub.subscribe(CHANNEL);

const receivers1 = await pub.publish(CHANNEL, "BBCA@9512");
// receivers1 is 1, and `received` now contains "BBCA@9512"

State subscription yang menentukan hasilnya. Pub/Sub nggak menyediakan acknowledgment pemrosesan atau replay untuk subscriber yang telat. Jumlah receiver-nya nggak membuktikan bahwa sebuah client udah menyelesaikan pekerjaannya.

Streams: pembaca yang telat dapat entry yang disimpan

Lalu saya menjalankan tes dengan bentuk yang sama terhadap sebuah stream. Tulis sekumpulan tick dulu, lalu baca lagi dengan client yang nggak hadir di satu pun proses penulisannya:

typescript
const STREAM = "demo:s04:stream";
const sim = new MarketSimulator({ seed: 7, volatilityScale: 4 });

for (const tick of sim.batch(12)) {
  await pub.xadd(STREAM, "*", ...tickToRedis(tick));
}

// A reader that arrives only now, after every write already happened:
const history = await pub.xrange(STREAM, "-", "+");
console.log(`a consumer that arrived AFTER the writes still sees all ${history.length} entries`);
// history.length is 12. Every tick survived.

Command XRANGE key - + membaca entry yang disimpan dari ID terendah sampai tertinggi. Pembaca saya yang telat berhasil mengambil kedua belas entry. Consumer group juga bisa membaca secara independen, tapi ID awalnya menentukan history mana yang mereka terima. Setting retention dan persistence tetap membatasi recovery.

Harga yang kamu bayar saat terputus

Kedua bagian tes itu adalah satu eksperimen yang dijalankan dengan satu kondisi yang berubah: apakah udah ada yang mendengarkan. Ubah itu jadi subscriber yang offline lalu kembali lagi, dan gambarannya jadi seperti ini:

mermaid
sequenceDiagram
    participant P as Producer
    participant Pub as "Pub/Sub Channel"
    participant S as "Stream Key"
    participant Sub as Subscriber
    Sub--)Pub: disconnects
    P->>Pub: PUBLISH "BBCA@9500"
    Note over Pub: 0 receivers, message discarded
    P->>S: XADD "BBCA@9500"
    Note over S: entry appended, stays until trimmed
    Sub->>Pub: reconnects, SUBSCRIBE
    Note over Sub: nothing published during the gap is recoverable
    Sub->>S: XRANGE - +
    Note over Sub: replays every entry, including "BBCA@9500"

Subscriber Pub/Sub yang terputus bakal melewatkan pesan yang dikirim selama jeda itu. Pembaca stream bisa memulihkan entry yang disimpan lewat XRANGE atau XREADGROUP. Consumer group Redis Streams menjelaskan Pending Entries List. Daftar ini melacak delivery yang butuh acknowledgement. Consumer lain bisa mengambil alihnya kalau consumer aslinya berhenti. Recovery tetap butuh data entry-nya masih tersedia.

Retention dan persistence stream punya batasan masing-masing. MAXLEN bisa menghapus entry lama dari memory. Snapshot RDB dan setting AOF menentukan apa yang bisa dipulihkan Redis dari disk. Pembaca yang terputus nggak bakal menghapus entry dengan sendirinya, tapi trimming, penghapusan, atau kegagalan bisa bikin entry jadi nggak tersedia.

Dua bentuk fan out, berdampingan

mermaid
graph TD
    subgraph "Pub/Sub - fire and forget"
        A["PUBLISH tick"] --> B{"Subscriber connected right now?"}
        B -->|No| C["Message discarded, 0 receivers"]
        B -->|Yes| D["Delivered once, no history kept"]
    end
    subgraph "Streams - durable log"
        E["XADD tick"] --> F["Entry appended to the log"]
        F --> G["Late reader: XRANGE replays full history"]
        F --> H["Consumer group: XREADGROUP + PEL + XACK"]
    end

Kapan Pub/Sub pas banget

Pakai Pub/Sub kalau update yang belakangan menggantikan yang sebelumnya dan update yang terlewat masih bisa diterima. Contohnya posisi cursor, typing indicator, dan permintaan ke client yang terhubung untuk me-refresh state terkini. Aplikasinya tetap harus mengambil state terkini setelah reconnect.

Layer UI dari sebuah price ticker yang mem-broadcast tick terbaru ke siapa pun yang lagi buka halamannya adalah use case Pub/Sub yang masuk akal, persis karena alasan itu. Client yang melewatkan tiga harga di antaranya waktu reconnect nggak perlu memutar ulang harga-harga itu. Yang dia butuhkan adalah harga yang sekarang. Harga dari kemudahan itu adalah Pub/Sub nggak bisa menjawab "what happened while I was gone," karena dia memang nggak pernah didesain untuk menyimpan jawabannya.

Kapan dia menghilangkan data yang kamu butuhkan

Kegagalannya datang dalam bentuk pesan yang terlewat, bukan crash atau exception, dan itulah yang bikin berbahaya. Sebuah worker bisa restart waktu deploy dan melewatkan setiap pesan yang di-publish selama dua detik dia mati. Client di balik koneksi mobile yang nggak stabil bisa terputus sebentar, persis di jendela waktu sebuah price alert terkirim. Nggak ada apa pun di Pub/Sub yang memberi tahu kamu kalau semua itu terjadi. PUBLISH mengembalikan sebuah angka, angkanya lebih kecil dari yang kamu kira, dan sistemnya jalan terus.

Kalau pekerjaan harus tetap selamat waktu consumer terputus, pakai storage yang menyimpan data dan sebuah proses recovery. Streams atau Kafka bisa mendukung ini. Tapi nggak satu pun dari keduanya bisa bikin efek bisnis eksternal tereksekusi tepat satu kali tanpa koordinasi di level aplikasi.

Pesan mana yang selamat

  • PUBLISH menyebarkan pesan cuma ke client yang subscribe tepat di saat itu. Nggak ada subscriber berarti pesannya dibuang, tanpa error dan tanpa jejak.
  • XADD menambahkan entry ke sebuah stream. Pembaca yang telat bisa mengambil entry yang disimpan dengan XRANGE. Setiap consumer group melacak posisinya sendiri dan delivery yang masih pending lewat XREADGROUP.
  • Panggilan publish saya mengembalikan 0 receiver waktu nggak ada subscriber yang terhubung. Pembaca stream yang telat berhasil memulihkan semua 12 entry dari tes itu.
  • MAXLEN membatasi retention stream. Recovery dari disk juga bergantung pada konfigurasi RDB atau AOF.
  • Pakai Pub/Sub untuk sinyal yang sekali pakai dan saling menggantikan: presence, cursor, "refetch now." Pakai Streams untuk apa pun yang harus selamat dari disconnect atau harus diproses dengan andal.

Ticker saya menaruh price tick BBCA dan alert di nomor partition yang sama di dua topic Kafka. Dengan assignment consumer yang cocok dan state yang tersedia, dia bisa melakukan join atas record-record ini di dalam satu proses.

TOPIK

SELANJUTNYA DI SERI INICo-Partitioning: Key Sama, Partition Sama, di Kedua Topic

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.