Kembali ke jurnalCATATAN FAJAR
Software Engineering5 menit baca

Co-Partitioning: Key Sama, Partition Sama, di Kedua Topic

Key, jumlah partition, dan partitioner yang sama menyelaraskan record antar topic, sementara join lokal juga butuh consumer assignment yang cocok.

BAGIAN 10 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 MAXLEN
  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 TopicKamu di sini
  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

"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.

typescript
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:

mermaid
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:

mermaid
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:

typescript
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.

mermaid
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 dan partitionForKey("BBCA", 8) adalah 2.
  • Redis Cluster memakai hash tag untuk menaruh key yang saling berkaitan di slot yang sama. Key market:{BBCA} dan notif:{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.

TOPIK

SELANJUTNYA DI SERI INIMenulis Custom Partition Assigner KafkaJS untuk Local Join

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.