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:
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:
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.XREADmengikuti ujung stream, kayaktail -f, sambil menunggu entry baru setelah ID tertentu.XLENmengembalikan jumlah entry.
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.
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.
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.
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.
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
MAXLENmembatasi 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.

Memuat komentar...