Kafka menjamin urutan di dalam setiap partition. Sebuah topic bisa punya urutan yang berbeda di antara partition-partition-nya, dan ini penting karena tick BBCA yang urutannya kacau bakal bikin ticker menampilkan harga terakhir yang salah.
Waktu saya membangun ticker dengan Kafka dan Redis Streams, hal pertama yang saya cek adalah routing record. Saya perlu memahami topic, partition, offset, dan key.
Log, dipotong-potong jadi partition
Sebuah topic di Kafka adalah log yang punya nama, dalam kasus saya market-updates. Topic dipecah jadi beberapa partition, dan setiap partition adalah rangkaian record miliknya sendiri yang berurutan, immutable, dan append only. Posisi sebuah record di dalam partition-nya adalah offset-nya, sebuah angka yang cuma bisa naik.
Kafka menjaga urutan append di dalam setiap partition. Kafka nggak menetapkan urutan antara record 100 di partition 0 dan record 50 di partition 3. Partition yang lebih banyak memungkinkan lebih banyak consumer di sebuah group membaca secara independen. Ini menaikkan paralelisme tanpa menciptakan urutan bersama di antara partition.
Key sebuah record menentukan partition-nya pada saat produce.
Key yang menentukan partition
Waktu producer mengirim sebuah record, dia bisa menyertakan key. Kalau key-nya ada, Kafka nggak memilih partition secara acak atau round robin, Kafka menghitungnya secara deterministik:
partition = (murmur2(key) & 0x7fffffff) % numberOfPartitions
Key ter-encode yang sama dipetakan ke partition yang sama selama jumlah partition dan partitioner-nya nggak berubah. Saya memakai simbol sebagai key. Jadi semua tick BBCA masuk ke satu partition sesuai urutan append di broker.
Ini nggak menetapkan urutan event time di antara producer yang independen. Kalau key-nya nggak ada, KafkaJS bisa menyebar simbol itu ke beberapa partition, dan cakupan urutan bersama ini pun hilang.
Codec saya yang pakai key cuma satu function. Sebuah Tick punya symbol, price, dan seq (counter per simbol yang saya pakai semata-mata untuk mendeteksi pelanggaran urutan nanti):
export function tickToKafka(tick: Tick): { key: string; value: string } {
return { key: tick.symbol, value: JSON.stringify(tick) };
}
Key-nya adalah string simbol, nggak ada yang lebih canggih dari itu, dan itulah yang jadi dasar routing partition.
graph LR
BBCA --> P0["Partition 0"]
BBRI --> P0
BMRI --> P0
BBNI --> P0
ANTM --> P1["Partition 1"]
ASII --> P2["Partition 2"]
GOTO --> P2
ICBP --> P2
UNVR --> P3["Partition 3"]
TLKM --> P5["Partition 5"]
Mengimplementasikan ulang murmur2 di JavaScript
KafkaJS udah menghitung hash-nya pada saat produce. Saya juga mengimplementasikan murmur2 di JavaScript untuk simulasi di browser dan test harness. Function itu memprediksi partition tanpa broker. Lalu saya bisa membandingkan hasilnya dengan partition yang dipakai KafkaJS.
const SEED = 0x9747b28c;
const M = 0x5bd1e995;
const R = 24;
export function murmur2(key) {
const data = typeof key === "string" ? new TextEncoder().encode(key) : key;
const length = data.length;
let h = SEED ^ length;
const blocks = Math.floor(length / 4);
for (let i = 0; i < blocks; i++) {
const i4 = i * 4;
let k =
(data[i4] & 0xff) |
((data[i4 + 1] & 0xff) << 8) |
((data[i4 + 2] & 0xff) << 16) |
((data[i4 + 3] & 0xff) << 24);
k = Math.imul(k, M);
k ^= k >>> R;
k = Math.imul(k, M);
h = Math.imul(h, M);
h ^= k;
}
// These cases intentionally fall through, matching Kafka's Java client.
const tail = blocks * 4;
switch (length - tail) {
case 3:
h ^= (data[tail + 2] & 0xff) << 16;
// Deliberate fall through to fold the next byte into the same mix.
case 2:
h ^= (data[tail + 1] & 0xff) << 8;
// Deliberate fall through to fold the final byte into the same mix.
case 1:
h ^= data[tail] & 0xff;
h = Math.imul(h, M);
}
h ^= h >>> 13;
h = Math.imul(h, M);
h ^= h >>> 15;
return h | 0;
}
export function partitionForKey(key, numPartitions) {
const positiveHash = murmur2(key) & 0x7fffffff;
return positiveHash % numPartitions;
}
Bagian yang ribet adalah Math.imul. Angka di JavaScript itu float, dan murmur2 bergantung pada perkalian signed 32 bit integer yang overflow dengan cara yang sama seperti di Java atau C. Math.imul memberi kamu perilaku wraparound itu. TextEncoder mengubah string key jadi byte UTF 8 sebelum hash dijalankan, sedangkan Uint8Array memungkinkan pemanggil memberikan byte yang persis secara langsung. Tail switch dan final mix-nya mengikuti Utils.murmur2(byte[]) dari Apache Kafka 4.3.1, dan partitionForKey menerapkan positive hash mask Kafka sebelum mengambil modulus.
Saya mengecek function yang diekstrak ini terhadap implementasi Python independen dari kontrak level byte yang sama. Vector-vector tetap ini menguji sisa satu, dua, dan tiga byte, ditambah sebuah key UTF 8 multibyte. Hash-nya berupa nilai signed 32 bit, dan partition-nya dihitung dengan enam partition:
| Input | Byte UTF 8 | Hash murmur2 | Partition |
|---|---|---|---|
a | 1, tail 1 | -1563381124 | 4 |
ab | 2, tail 2 | 316155434 | 2 |
abc | 3, tail 3 | 479470107 | 3 |
abcd | 4 | -1323649548 | 2 |
é | 2, tail 2 | 186971271 | 3 |
Enam partition, satu hot spot
Saya menyiapkan market-updates dengan 6 partition, paralelisme yang cukup untuk 10 simbol tanpa setiap simbol harus dapat partition sendiri. Menjalankan partitionForKey atas daftar simbol menunjukkan dengan tepat ke mana masing-masing simbol pergi:
| Simbol | Partition |
|---|---|
| BBCA | 0 |
| BBRI | 0 |
| BMRI | 0 |
| BBNI | 0 |
| ANTM | 1 |
| ASII | 2 |
| GOTO | 2 |
| ICBP | 2 |
| UNVR | 3 |
| TLKM | 5 |
Empat dari sepuluh simbol dipetakan ke partition 0, sementara partition 4 nggak kebagian satu pun. Set key yang kecil ini menghasilkan distribusi yang nggak merata. Simbol yang ramai bakal menambah beban ke partition tempat dia di-assign.
Partition yang lebih banyak bisa mendistribusikan ulang key, tapi nggak bisa memecah satu key secara otomatis. Composite key atau custom partitioner bisa mengubah distribusi beban, dengan konsekuensi cakupan urutan yang berbeda. Cek pemetaannya terhadap traffic yang diperkirakan sebelum memilih strategi.
Membuktikannya di live broker
Baca dokumentasi itu satu hal, percaya sama dokumentasi itu di depan live broker itu lain cerita. Jadi saya menjalankan pertanyaan ini end to end terhadap Kafka di Docker, dalam empat langkah.
- Prediksi dulu: Panggil
partitionForKeyuntuk ke-10 simbol. Cetak partition yang diharapkan untuk setiap simbol. Perhitungan ini nggak menghubungi broker. - Produce dengan key: generate 60 tick simulasi, kirim ke
market-updates, masing-masing dengan keytick.symbol. - Consume dan cek: Baca ke-60 record. Bandingkan setiap partition dengan prediksinya. Untuk setiap simbol, cek bahwa
seqnaik secara strict. Nilai yang turun menandakan ada error urutan. - Lalu sengaja dirusak: Ulangi simulasinya dengan topic kedua. Kirim tick yang sama tanpa key. Sekarang tick setiap simbol tersebar ke beberapa partition, dan nggak ada consumer yang bisa melihatnya secara berurutan lagi.
sequenceDiagram
participant Prod as Producer
participant B as Broker "partition 0"
participant Cons as Consumer
Prod->>B: send "key=BBCA, seq=1"
Prod->>B: send "key=BBCA, seq=2"
Prod->>B: send "key=BBCA, seq=3"
Note over B: same key, same partition, every time
B->>Cons: seq=1
B->>Cons: seq=2
B->>Cons: seq=3
Note over Cons: seq strictly increasing, ordering held
Di run yang pakai key, setiap simbol dipetakan ke satu partition dan seq nggak pernah turun. Di run tanpa key, simbol yang tick-nya cukup banyak muncul di dua partition atau lebih. Tesnya melaporkan distribusi itu secara otomatis.
Kembali ke ticker
Saat ini ticker-nya melacak 10 simbol dan mungkin butuh lebih banyak. Partition memungkinkan consumer yang terpisah memproses kelompok simbol yang berbeda. Key simbol yang stabil menjaga simbol itu tetap di satu partition, tempat Kafka menjaga urutan append. Producer dan handler-nya juga harus menjaga urutan aplikasi apa pun yang dibutuhkan.
Ini juga menyiapkan sesuatu yang bakal saya bahas lagi nanti. notifications punya 6 partition yang sama dan key symbol yang sama, bukan userId. Ini disengaja: tick BBCA dan alert BBCA di-hash ke nomor partition yang sama, yang jadi penting begitu kamu mulai men-join kedua stream itu.
Aturan routing dalam sekali lihat
Kafka menjaga urutan di dalam setiap partition. Ekspresi (murmur2(key) & 0x7fffffff) % numPartitions memetakan key yang ter-encode ke sebuah partition. Function independen saya memakai Math.imul untuk perkalian 32 bit. Dengan function itu saya bisa membandingkan routing yang diprediksi dengan hasil dari broker.
Set key yang kecil bisa menghasilkan distribusi yang nggak merata. Sepuluh simbol saya di enam partition menaruh empat simbol di satu partition dan nggak ada satu pun di partition lain. Tesnya bikin ketimpangan ini kelihatan sebelum deployment.
Redis Streams nggak punya partition untuk dibagikan. Urutan dijaga dengan cara yang berbeda, yang saya bahas berikutnya di XADD, entry ID, dan MAXLEN, sekalian di titik mana cerita itu berhenti mirip dengan Kafka.

Memuat komentar...