Kembali ke jurnalCATATAN FAJAR
Software Engineering5 menit baca

Dasar-Dasar Redis Streams: XADD, Entry ID, dan MAXLEN

Pakai XADD, entry ID, stream key, dan MAXLEN untuk menulis, membaca, dan menyimpan data stream Redis.

BAGIAN 3 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 MAXLENKamu di sini
  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 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

Redis Streams, yang diperkenalkan di Redis 5.0, menyimpan rangkaian entry yang terurut di dalam sebuah key Redis. Tipe data ini mendukung history yang disimpan dan consumer group. Model urutan dan retention-nya berbeda dari partition Kafka.

Saya membandingkan kedua sistem ini sambil membangun ticker yang sama dua kali, dengan broker yang jalan di Docker di laptop saya. Ticker-nya mem-publish update harga dan mengevaluasi alert. Post ini menjelaskan entry ID, pembacaan, dan retention di Redis.

Log di dalam cache

Tambahkan entry dengan XADD. Redis bisa memberikan entry ID secara otomatis. Setiap entry berisi pasangan field dan value, yang di sini direpresentasikan sebagai string:

text
XADD market-updates * symbol BBCA price 9530 seq 1
=> "1718000000123-0"

Tanda * bilang ke Redis "give me the next ID." Tick dari market simulator saya masuk dengan cara yang sama, lewat wrapper ioredis yang kecil:

typescript
import { createRedis } from "../../lib/redis/client";
import { MarketSimulator, tickToRedis } from "../../lib/domain";

const redis = createRedis();
const sim = new MarketSimulator({ seed: 3, volatilityScale: 5 });

for (const tick of sim.batch(8)) {
  // '*' tells Redis to assign the next time-ordered ID.
  const id = await redis.xadd(KEY, "*", ...tickToRedis(tick));
}

tickToRedis mengubah sebuah tick jadi pasangan field dan value yang diterima XADD: symbol, price, prevPrice, volume, ts, dan seq. Arti field-field ini ditentukan oleh aplikasi.

Entry ID bukan offset

Setiap entry ID punya format <millisecondsTime>-<sequence>, dan nilainya selalu naik.

Offset Kafka menandai sebuah posisi di dalam partition. Timestamp record-nya terpisah. ID yang dibuat Redis menggabungkan komponen millisecond dengan komponen sequence.

Kalau dua entry memakai komponen millisecond yang sama, Redis menaikkan sequence-nya, misalnya ...123-0 dan ...123-1. Ketika waktu bergerak maju, sequence-nya bisa mulai lagi dari 0. Kalau jam bergerak mundur, Redis tetap menjaga ID yang naik dengan memakai komponen waktu sebelumnya. ID yang ditentukan secara eksplisit juga bisa berbeda dari wall clock time.

Entry ID Redis adalah cursor ke dalam stream. Untuk ID yang dibuat otomatis, komponen millisecond-nya juga mendukung query rentang waktu dengan XRANGE. Jangan anggap setiap entry ID yang mungkin ada sebagai timestamp kedatangan yang persis.

Membacanya lagi: XRANGE dan XREAD

Tiga command ini mencakup sebagian besar yang kamu butuhkan:

  • XRANGE key - + membaca sebuah rentang berdasarkan ID (- dan + berarti seluruh stream). Cocok untuk replay dan inspeksi.
  • XREAD mengikuti ujung stream, kayak tail -f, sambil menunggu entry baru setelah ID tertentu.
  • XLEN mengembalikan jumlah entry.
typescript
const range = await redis.xrange(KEY, "-", "+");
for (const [id, flat] of range.slice(0, 4)) {
  const f: Record<string, string> = {};
  for (let i = 0; i + 1 < flat.length; i += 2) f[flat[i]] = flat[i + 1];
  console.log(`${id}  ${f.symbol} price=${f.price} seq=${f.seq}`);
}

Redis mengembalikan entry dalam bentuk pasangan [id, flatFieldArray]. Aplikasi mengubah setiap array field jadi sebuah object. Command-nya nggak menyediakan schema registry.

Satu stream, satu partition

Satu stream key adalah satu log yang terurut total, padanan langsung dari satu partition Kafka. Redis nggak punya konsep bawaan untuk mempartisi satu key.

Key routing di Kafka memilih dari sekumpulan partition yang jumlahnya tetap. Topic market-updates saya punya enam. Topic itu merutekan simbol dengan (murmur2(key) & 0x7fffffff) % numberOfPartitions.

Sepuluh simbol harus berbagi enam partition itu. Empat simbol saya memilih partition 0. Record untuk setiap simbol tetap berkumpul, tapi distribusinya nggak membagi beban secara merata.

Untuk membagi data stream Redis, aplikasinya membuat beberapa stream key. Misalnya, market:{BBCA} dan market:{BBRI} masing-masing menampung satu simbol. Setiap stream menjaga urutan entry-nya sendiri. Desain ini mengelola satu stream per simbol, bukan sejumlah partition Kafka yang tetap. Redis Cluster tetap meng-hash stream key itu ke dalam slot.

mermaid
graph TD
    A["market-updates (Kafka topic)"] --> B["6 partitions, murmur2(key) picks one"]
    B --> C["10 symbols hash into 6 buckets, collisions guaranteed"]
    D["Redis: no topic, just keys"] --> E["market:{BBCA}"]
    D --> F["market:{BBRI}"]
    D --> G["market:{...}, one stream per symbol"]

Jumlah partition yang tetap membatasi jumlah log yang harus dikelola, tapi beberapa key bisa berbagi partition. Satu stream per simbol memisahkan entry mereka ke log yang berbeda. Jumlah stream-nya lalu ikut bertambah seiring jumlah simbol. Redis Cluster bisa mendistribusikan stream-stream ini, walaupun key yang berbeda tetap bisa berbagi slot atau node.

Tinggalnya di RAM, jadi kamu yang men-trim-nya

Topic Kafka menyimpan log segment di disk dan menerapkan aturan retention yang dikonfigurasi. Redis menyimpan entry stream di memory. Tanpa trimming atau policy penghapusan lain, stream-nya terus membesar seiring entry yang datang.

typescript
await redis.xadd(KEY, "MAXLEN", "~", "10", "*", ...tickToRedis(tick));

MAXLEN ~ N mengatur batas panjang secara kira-kira. Trimming persis dengan MAXLEN N menghapus entry secukupnya supaya batasnya terpenuhi dengan tepat. Trimming kira-kira menghapus macro node internal secara utuh, yang bisa mengurangi kerja. Karena itu, stream-nya bisa menyimpan lebih dari N entry.

mermaid
graph LR
    P["Producer: XADD"] --> S["Stream in RAM"]
    S --> M{"MAXLEN set?"}
    M -->|"No"| G["Memory grows unbounded"]
    M -->|"Yes, MAXLEN ~ N"| T["Approximate trim, cheap"]

Persistence bergantung pada konfigurasi. Recovery dari RDB bisa kehilangan write sejak snapshot terakhir. Dengan AOF dan appendfsync everysec, sebuah crash bisa menghilangkan sekitar satu detik write. Policy always menyinkronkan write lebih sering dengan biaya tambahan.

Container saya pakai --appendonly yes --appendfsync everysec. Bandingkan setting persistence dan replication sebelum bikin klaim soal durability, baik untuk Redis maupun Kafka.

Panggilan client dan koneksi

Setup client-nya memakai dua connection factory Redis yang terpisah.

typescript
export function createRedis(opts: RedisOptions = {}): Redis {
  return new Redis(env.REDIS_URL, opts);
}

export function createBlockingRedis(opts: RedisOptions = {}): Redis {
  return new Redis(env.REDIS_URL, { maxRetriesPerRequest: null, ...opts });
}

Factory createRedis menangani command biasa seperti XADD, XLEN, dan pipeline. Factory createBlockingRedis menyediakan koneksi terpisah untuk XREAD atau XREADGROUP dengan BLOCK. Command lain di koneksi yang sedang ter-block harus menunggu.

Setting maxRetriesPerRequest: null menghapus batas retry per request di ioredis. Setting ini mengatur perilaku retry, bukan blocking timeout di Redis. Waktu tunggu itu diatur oleh argumen BLOCK.

Ringkasan model Redis

  • Redis Stream menyimpan rangkaian entry yang terurut di dalam sebuah key.
  • Entry ID memakai format <millis>-<seq> dan nilainya naik di dalam stream. ID yang dibuat otomatis biasanya mencerminkan waktu server, dengan tetap mengikuti aturan monotonic.
  • Satu stream nggak punya partitioning internal. Stream yang lebih banyak memungkinkan routing dan worker yang terpisah.
  • Opsi MAXLEN membatasi entry yang disimpan. Trimming kira-kira dengan ~ bisa menyisakan entry lebih banyak dari target.
  • Setting AOF, RDB, replication, dan retention menentukan batas recovery.

Pertanyaan berikutnya adalah apa yang terjadi kalau consumer yang lagi memegang satu batch entry tumbang di tengah jalan. Saya bakal mulai dari jawaban Kafka, lalu beralih ke Redis: consumer group, lag, dan apa yang dilakukan rebalance terhadap pekerjaan yang lagi berjalan.

TOPIK

SELANJUTNYA DI SERI INIKafka Consumer Group: Lag, Draining, dan Rebalancing

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.