Awalnya ticker ini cuma pakai satu broker Kafka dan satu instance Redis. Setup itu cukup untuk mengetes consumer group dan PEL, tapi nggak bisa untuk mengetes node yang mati. Saya menambahkan cluster dengan tiga broker Kafka dan Redis Cluster dengan tiga node. Keduanya jalan di file Docker Compose yang sama.
Ketika satu node nggak cukup
Profile compose kedua menjalankan enam container: kvr-kc1/2/3 untuk Kafka (mode KRaft, tanpa ZooKeeper, port 19092/19094/19096), dan kvr-rc1/2/3 untuk Redis (cluster mode aktif, port 7000/7001/7002). Dua cluster, bentuknya sama di atas kertas: masing-masing tiga node. Yang membedakan kedua sistem ini adalah apa yang mereka lakukan dengan tiga node itu.
Di tes ini, Kafka menyimpan salinan setiap partition di tiga broker. Redis membagi key ke tiga node primary tanpa replica. Konfigurasi Kafka-nya mengetes replication dan failover. Konfigurasi Redis-nya mengetes sharding. Kedua produk ini bisa memakai replication dan distribusi data sekaligus.
Replication di Kafka: leader, follower, dan ISR
Konfigurasi broker-nya memakai KAFKA_DEFAULT_REPLICATION_FACTOR: "3" dan KAFKA_MIN_INSYNC_REPLICAS: "2". Log internal untuk offset dan transaction juga memakai replication factor 3. Setiap partition punya satu leader dan dua follower di broker yang berbeda. Di setup ini, client memakai leader, dan follower menyalin record dari leader.
ISR (in sync replica set) berisi replica yang memenuhi syarat sinkronisasi Kafka. Partition yang sehat punya semua replica-nya di dalam ISR. Kafka mengeluarkan replica yang tertinggal terlalu jauh atau yang offline. Read butuh leader yang tersedia. Write juga bergantung pada setting acknowledgement dan ukuran minimum ISR.
Test harness saya membaca state persis seperti ini dari admin client:
interface PartInfo {
partition: number;
leader: number;
replicas: number[];
isr: number[];
}
async function partitions(admin: Admin): Promise<PartInfo[]> {
const meta = await admin.fetchTopicMetadata({ topics: [TOPIC] });
return meta.topics[0].partitions
.map((p) => ({ partition: p.partitionId, leader: p.leader, replicas: p.replicas, isr: p.isr }))
.sort((a, b) => a.partition - b.partition);
}
Topic replicationFactor: 3 yang sehat punya tiga replica di ISR setiap partition. Dengan acks=all, leader menunggu ISR saat ini meng-acknowledge write tersebut. Setting min.insync.replicas menentukan ukuran minimum ISR yang dibutuhkan supaya write itu diterima.
Dengan min.insync.replicas=2 dan RF=3, write tetap bisa jalan setelah satu broker mati. Kalau cuma tersisa satu replica di ISR, write ini gagal.
Mematikan broker dan melihat failover terjadi
Saya memaksa satu broker gagal setelah mem-produce 30 record yang punya key. Saya pakai idempotent producer dari contoh transaction. Tesnya mencari leader dari partition 0 lalu menghentikan container-nya:
const victim = before[0].leader;
execSync(`docker stop kvr-kc${victim}`, { stdio: "ignore" });
let after = before;
for (let i = 0; i < 15; i++) {
await sleep(2000);
try {
after = await partitions(admin);
} catch {
continue; // metadata fetch may blip while the broker drops
}
if (after[0].leader !== victim) break;
}
Tesnya melakukan polling fetchTopicMetadata setiap dua detik sampai partition 0 punya leader baru. Controller memilih salah satu anggota ISR yang tersisa. Di failover yang bersih ini, ISR turun dari tiga anggota jadi dua. Record berikutnya dengan acks=all berhasil karena ISR masih memenuhi min.insync.replicas=2.
sequenceDiagram
participant Admin
participant B1 as Broker 1
participant B2 as Broker 2
participant B3 as Broker 3
Admin->>B1: describeCluster shows leader p0 is broker 1, ISR is 1 2 3
Note over B1: docker stop kvr-kc1
Admin->>B2: fetchTopicMetadata polling
Note over B2,B3: controller elects new leader from the remaining ISR
Admin->>B2: leader p0 is now broker 2, ISR is 2 3
Note over Admin: producer.send still succeeds with acks=all since ISR 2 meets min.insync.replicas 2
Note over B1: docker start kvr-kc1
B1->>B2: rejoin and replicate to catch up
Admin->>B2: leader p0 is broker 2, ISR healed back to 1 2 3
Lalu saya menyalakan ulang container-nya dan polling sampai setiap partition punya tiga replica di ISR lagi. Ini memperlihatkan broker yang baru di-restart mengejar ketertinggalannya. Dengan acks=all, min.insync.replicas menentukan ukuran minimum ISR yang dibutuhkan supaya write berhasil.
Redis Cluster: sharding keyspace ke 16384 slot
Redis Cluster ini punya tiga node master. Masing-masing memegang rentang yang terpisah dari 16384 hash slot. Saya pakai redis-cli --cluster create dengan tiga announce IP dan tanpa replica. Setup ini memperlihatkan sharding, tanpa replica yang bisa menggantikan master yang mati.
Setiap key di-hash ke sebuah slot dengan CRC16, dan slot itu yang menentukan node mana yang memilikinya:
/** CRC16/XMODEM, exactly as Redis computes it (poly 0x1021, init 0). */
export function crc16(input: string): number {
const bytes = new TextEncoder().encode(input);
let crc = 0;
for (const byte of bytes) {
crc ^= byte << 8;
for (let i = 0; i < 8; i++) {
crc = (crc & 0x8000) ? ((crc << 1) ^ 0x1021) & 0xffff : (crc << 1) & 0xffff;
}
}
return crc & 0xffff;
}
/** The cluster slot for a key, honoring hash tags {...}. */
export function keySlot(key: string): number {
const open = key.indexOf("{");
if (open !== -1) {
const close = key.indexOf("}", open + 1);
if (close > open + 1) key = key.slice(open + 1, close);
}
return crc16(key) % 16384;
}
Client Cluster dari ioredis memakai CLUSTER SLOTS untuk membaca kepemilikan slot. Lalu saya memanggil keySlot untuk masing-masing dari sepuluh symbol. Container-nya mengiklankan alamat Docker internal (172.30.0.1x:6379) yang nggak bisa dijangkau dari host.
Saya menghabiskan sebagian malam untuk masalah koneksi ini. natMap memetakan setiap alamat internal ke port localhost yang di-publish. Tanpa itu, koneksi pertama berhasil tapi request yang di-redirect gagal.
Hash tag: jawaban Redis untuk co partitioning
Operasi atomic pada beberapa key di Redis Cluster mengharuskan key-key itu berada di slot yang sama. Tanpa hash tag yang sama, market:{BBCA} dan notif:{BBCA} bisa jatuh ke slot yang berbeda. Redis cuma meng-hash substring di dalam {...} kalau key-nya punya hash tag yang valid. Jadi kedua key itu berbagi slot dan node yang sama. Lua script atau transaction MULTI bisa mengakses keduanya sekaligus.
await cluster.xadd("market:{BBCA}", "*", "p", "9500");
await cluster.xadd("notif:{BBCA}", "*", "p", "alert");
const lua = "return { redis.call('XLEN', KEYS[1]), redis.call('XLEN', KEYS[2]) }";
const ok = await cluster.eval(lua, 2, "market:{BBCA}", "notif:{BBCA}");
// -> one node, one atomic call
Kalau tag-nya disilangkan, cluster langsung menolak:
try {
await cluster.eval(lua, 2, "market:{BBCA}", "notif:{BBRI}");
} catch (e) {
// rejected: CROSSSLOT Keys in request don't hash to the same slot
}
Key market:{BBCA} dan notif:{BBRI} memakai hash tag yang berbeda dan bisa dipetakan ke slot yang berbeda. Script yang butuh kedua key ada di satu slot pun gagal dengan CROSSSLOT.
Custom assigner KafkaJS menyelesaikan masalah yang mirip, yaitu penempatan proses. Assigner itu memberikan partition yang cocok dari beberapa topic ke satu consumer pod. Sebaliknya, hash tag di Redis mengatur penempatan data di sebuah slot. Keduanya adalah mekanisme terpisah dengan jaminan yang berbeda.
graph TD
subgraph Slots["16384 hash slots split across 3 masters"]
N1[rc1 master]
N2[rc2 master]
N3[rc3 master]
end
K1["market:{BBCA}"] -->|hash tag BBCA| N2
K2["notif:{BBCA}"] -->|hash tag BBCA| N2
K3["notif:{BBRI}"] -->|hash tag BBRI| N3
Docker Compose yang sama, dua alasan berbeda untuk menambah node
Kalau kedua cluster ini ditaruh berdampingan, filosofinya hampir berlawanan:
| Kafka (cluster RF=3) | Redis Cluster (3 master) | |
|---|---|---|
| Isi setiap node | Salinan penuh dari partition yang di-assign ke node itu | Potongan terpisah dari 16384 slot |
| Kenapa kamu menambah node | Durability dan availability | Kapasitas dan throughput |
| Apa yang terjadi kalau satu node mati | Leader di-failover ke replica yang in sync, nggak ada data yang udah di-acknowledge yang hilang | Slot di node itu nggak tersedia sampai ada replica yang mengambil alih, atau hilang kalau node itu nggak punya replica |
| Mekanisme pengamannya | ISR ditambah min.insync.replicas ditambah acks=all | Replica per shard (nggak dikonfigurasi di setup saya) |
| Masalah koordinasi multi key | Co-partitioning: taruh key yang saling terkait di partition yang sama | Hash tag: taruh key yang saling terkait di slot yang sama |
Replication di Kafka memungkinkan replica lain melayani sebuah partition setelah leader-nya gagal. Sharding di Redis mendistribusikan key ke beberapa node untuk menambah kapasitas. Saya mengonfigurasi tes ini secara terpisah supaya bisa mengamati masing-masing perilaku. Redis Cluster yang punya replica bisa menyediakan sharding dan failover sekaligus, dengan tetap tunduk pada batasan persistence dan replication-nya.
Replication dan sharding dalam satu gambaran
- Replication factor di Kafka menentukan jumlah replica dari sebuah partition. ISR berisi replica yang memenuhi syarat sinkronisasinya.
acks=allmenunggu ISR saat ini.min.insync.replicasmengatur ukuran minimum ISR yang dibutuhkan supaya write berhasil.- Di tes RF=3 ini, anggota ISR yang lain menggantikan leader yang dihentikan. Write jalan lagi kalau ISR yang tersedia memenuhi batas minimum yang dikonfigurasi. Broker yang di-restart mengejar ketertinggalannya dan bergabung lagi ke ISR.
- Redis Cluster membagi 16384 slot yang jumlahnya tetap ke node-node master memakai CRC16 dari key, atau dari hash tag di dalam
{...}kalau ada. - Hash tag menaruh key yang saling terkait di satu slot. Lua script bisa mengakses key-key itu secara atomic. Operasi yang butuh satu slot bakal gagal dengan
CROSSSLOTkalau key-nya tersebar di beberapa slot. - Sharding membagi data ke beberapa node. Replication membuat salinan. Redis Cluster bisa memakai keduanya, walaupun tes ini nggak punya replica.
Sekarang tesnya udah mencakup broker yang gagal. Post berikutnya menambahkan state ke ticker lewat window OHLC dan price alert yang merespons ketika harga melewati threshold atau berdasarkan level saat ini.

Memuat komentar...