"I want partition 3 of notifications consumed on the same pod as partition 3 of market-updates." Kalimat itu menggambarkan apa yang dibutuhkan sistem ticker kecil saya, dan di dalamnya tersembunyi dua masalah terpisah yang terus saya campuradukkan selama berminggu-minggu.
Kebutuhan pertama soal data: nomor partition yang sama harus berisi simbol yang sama. Kebutuhan kedua soal assignment: satu proses consumer harus memegang kedua partition yang bersesuaian itu. Hash function yang deterministik bisa memenuhi kebutuhan pertama.
Satu topic, satu key, satu partition
Post tentang ordering menjelaskan routing function-nya: (murmur2(key) & 0x7fffffff) % numPartitions. Selama byte key, partitioner, dan jumlah partition nggak berubah, key yang sama bakal memilih partition yang sama. market-updates memakai simbol sebagai key-nya supaya tick BBCA tetap berada di satu partition yang berurutan.
Itu pernyataan tentang satu topic. Co partitioning adalah apa yang terjadi saat kamu menerapkan resep yang persis sama ke dua topic sekaligus.
Definisi co partitioning
Dua topic disebut co partitioned kalau keduanya punya:
- jumlah partition yang sama, dan
- key yang sama, yang di-route oleh partitioner yang sama.
Di ticker yang saya bangun dua kali, market-updates dan notifications sama-sama punya 6 partition dan memakai symbol sebagai key. Function partitionForKey menerima key dan jumlah partition. Function ini nggak menerima nama topic. Untuk key ter-encode dan jumlah partition yang sama, function ini mengembalikan angka yang sama untuk kedua topic.
export function partitionForKey(key: string, numPartitions: number): number {
return (murmur2(key) & 0x7fffffff) % numPartitions;
}
Kafka nggak butuh setting khusus untuk properti ini. Kedua topic harus memakai jumlah partition, encoding key, dan logika partition deterministik yang sama.
Di mana setiap simbol mendarat, di kedua topic
Karena input-nya identik, tabel output-nya juga identik. Partition mana pun yang diberikan partitionForKey untuk BBCA di market-updates, BBCA juga dapat nomor yang sama di notifications:
graph LR
subgraph MU["market-updates (6 partitions)"]
MP0["Partition 0"]
MP1["Partition 1"]
MP2["Partition 2"]
MP3["Partition 3"]
MP5["Partition 5"]
end
subgraph NT["notifications (6 partitions)"]
NP0["Partition 0"]
NP1["Partition 1"]
NP2["Partition 2"]
NP3["Partition 3"]
NP5["Partition 5"]
end
BBCA --> MP0
BBCA --> NP0
ANTM --> MP1
ANTM --> NP1
ASII --> MP2
ASII --> NP2
UNVR --> MP3
UNVR --> NP3
TLKM --> MP5
TLKM --> NP5
BBCA dipetakan ke partition 0 di kedua topic, dan ANTM ke partition 1 di keduanya. Input yang sama ke hash function-lah yang menghasilkan ini. Nggak perlu lookup table atau link antar topic yang terpisah.
Apa yang merusaknya
Setiap bagian dari definisi itu penting, dan menghilangkan salah satunya bakal merusak co partitioning.
Jumlah partition berbeda: modulus adalah bagian dari formulanya. Perubahan jumlah bisa mengubah partition untuk sebuah key. Contohnya, BBCA dipetakan berbeda dengan 8 partition dibanding dengan 6:
graph TD
K["key: BBCA"] --> M6["mod 6, market-updates: partition 0"]
K --> M8["mod 8, notifications: partition 2"]
M6 -.->|"misaligned"| M8
Saya mengeceknya, bukan cuma berasumsi: partitionForKey("BBCA", 6) mengembalikan 0, dan partitionForKey("BBCA", 8) mengembalikan 2. Key sama, modulus beda, partition beda, jadi co partitioning udah hilang bahkan sebelum faktor lain ikut berperan.
Key berbeda: kalau notifications memakai userId, bukan symbol, user-lah yang menentukan partition. Alert untuk BBCA dari user yang berbeda-beda jadi bisa dipetakan ke partition yang berbeda dari tick BBCA.
Partitioner berbeda: kedua producer harus memakai logika hash dan encoding key yang sama. Message tanpa key atau custom partitioner yang berbeda bakal menghilangkan jaminan partition yang bersesuaian. Jumlah partition yang sama saja nggak cukup.
Co partitioning butuh aturan routing yang sama. Perubahan jumlah partition di kemudian hari bisa merusak keselarasan itu untuk record baru.
Padanannya di Redis: hash tag dan slot CRC16
Redis Cluster nggak punya partition dan nggak ada numPartitions yang kamu set per key. Setiap key di-hash ke salah satu dari 16384 slot tetap memakai CRC16 (varian XMODEM, polynomial 0x1021), lalu slot-slot itu didistribusikan ke berbagai node:
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;
}
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); // only the tag is hashed
}
}
return crc16(key) % 16384;
}
Key market:BBCA dan notif:BBCA biasanya di-hash secara independen karena string lengkapnya berbeda. Helper keySlot cuma meng-hash teks di dalam {...} kalau menemukan hash tag yang valid. Dengan market:{BBCA} dan notif:{BBCA}, input hash keduanya sama-sama "BBCA". Jadi, kedua key memakai slot dan node yang sama.
Saya menjalankan itu ke semua 10 simbol yang diperdagangkan ticker saya dan menghitung co location-nya. Dengan hash tag: 10 dari 10 pasangan simbol ada di slot yang sama, setiap kali, karena memang begitu desainnya. Tanpa hash tag: 0 dari 10, karena penggabungan string biasa nggak memberi CRC16 alasan untuk menghasilkan nilai yang sama.
graph LR
A["market:{BBCA}"] --> S1["slot"]
B["notif:{BBCA}"] --> S1
C["market:BBCA"] --> S2["slot"]
D["notif:BBCA"] --> S3["different slot"]
Key di slot yang sama bisa ikut dalam satu operasi MULTI/EXEC atau satu Lua script. Ini memungkinkan operasi Redis yang atomic pada stream yang saling berkaitan. Nomor partition di Kafka mengidentifikasi kelompok data, tapi nomor yang sama saja nggak menempatkan partition di broker atau consumer yang sama.
Layout bukan berarti penempatan
Jumlah partition dan pemetaan key yang sama menempatkan simbol yang sama di partition yang bersesuaian di market-updates dan notifications. Consumer assignment itu urusan terpisah. Assignment round robin di KafkaJS bisa menaruh partition yang bersesuaian di member yang berbeda untuk ukuran group tertentu.
Tahu bahwa sebuah simbol selalu dipetakan ke nomor partition yang sama di kedua topic itu perlu, tapi belum cukup. Runtime-nya tetap harus diberi tahu secara eksplisit untuk menaruh nomor partition yang bersesuaian di proses yang sama.
Apa yang dijamin co partitioning
- Kedua topic harus memakai jumlah partition, encoding key, dan partitioner deterministik yang sama. Dengan begitu, sebuah key dipetakan ke nomor partition yang sama.
partitionForKey(key, numPartitions)menerima input yang identik untuk kedua topic. Nggak perlu link antar topic yang terpisah.- Perubahan jumlah partition bisa mengubah pemetaan untuk record berikutnya. Record yang udah ada tetap di partition aslinya. Untuk BBCA,
partitionForKey("BBCA", 6)adalah 0 danpartitionForKey("BBCA", 8)adalah 2. - Redis Cluster memakai hash tag untuk menaruh key yang saling berkaitan di slot yang sama. Key
market:{BBCA}dannotif:{BBCA}memakai satu dari 16384 slot. Di tes 10 simbol saya, pasangan yang di-tag cocok 10 dari 10 kali dan pasangan biasa cocok 0 dari 10. - Join Kafka lokal juga butuh consumer assignment yang bersesuaian. Pemetaan partition saja nggak menentukan proses mana yang memegangnya.
Saya menulis custom partition assigner untuk KafkaJS supaya partition yang bersesuaian ditaruh di proses yang sama. Post berikutnya menjelaskan implementasinya.

Memuat komentar...