Kembali ke jurnalCATATAN FAJAR
Software Engineering6 menit baca

Menulis Custom Partition Assigner KafkaJS untuk Local Join

Memberikan partition yang cocok dari dua topic ke consumer KafkaJS yang sama, supaya consumer itu bisa join record-nya secara lokal.

BAGIAN 11 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 Topic
  11. 11Menulis Custom Partition Assigner KafkaJS untuk Local JoinKamu di sini
  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 8 bagian

Dengan jumlah partition, symbol key yang di-encode, dan partitioner yang sama, market-updates dan notifications memetakan BBCA ke nomor partition yang sama. Consumer assigner yang menentukan proses mana yang memiliki partition-partition itu. Nomor yang cocok aja nggak cukup untuk menaruh kedua partition di consumer yang sama.

Post soal co partitioning membiarkan masalah consumer assignment belum terjawab. Saya menghabiskan satu malam untuk mengimplementasikan assigner supaya ticker saya bisa join harga dan notifikasi di dalam setiap proses.

Co partitioning itu soal tata letak data, bukan penempatan

Kedua topic punya 6 partition dan pakai symbol key yang sama. Mapping default-nya menerapkan (murmur2(key) & 0x7fffffff) % numPartitions. Jadi, nomor partition yang sama berisi simbol yang sama. Consumer assignment harus menjaga keselarasan itu antar proses.

Itu jaminan di sisi data. Co partitioning memberi tahu di mana data berada, sedangkan consumer assignment menentukan proses mana yang membacanya. Bagian kedua ini namanya co localization, dan sepenuhnya bergantung pada cara kerja partition assignment di consumer group kamu.

Kenapa round robin assigner KafkaJS merusak ini

Java client Kafka punya RangeAssignor yang melakukan co localization untuk partition bernomor sama di semua topic yang di-subscribe sebuah consumer. Di JVM, assignment yang ramah untuk join bisa didapat hampir gratis.

Versi KafkaJS yang dipakai di sini punya round robin assigner. Assigner ini menyebarkan flat list berisi partition dari semua topic ke para member. Hasilnya bisa menyelaraskan nomor partition yang sama untuk beberapa kombinasi jumlah member dan jumlah partition.

Setup saya punya empat pod dan enam partition per topic. Assignment round robin-nya menaruh beberapa partition yang cocok di pod yang berbeda. Record-nya masih punya nomor partition yang sama, tapi local join butuh partition dari kedua topic ada di proses yang sama.

Saya buang waktu lebih banyak dari yang mau saya akui karena mengira perilaku di JVM juga berlaku di sini. Ternyata nggak. Perbaikannya harus dilakukan di layer assignment, di TypeScript, karena nggak ada assigner dari library yang bisa langsung dipakai.

Interface assigner di group protocol

KafkaJS menyediakan factory PartitionAssigner yang punya akses ke metadata cluster. Object yang di-return menyediakan name, version, protocol(), dan assign(). Member mengumumkan protokol yang mereka dukung waktu bergabung ke group. Group leader yang terpilih menjalankan assign() dan me-return assignment untuk dibagikan oleh broker.

typescript
import { AssignerProtocol, type PartitionAssigner } from "kafkajs";

const NAME = "CoPartitionAssigner";
const VERSION = 1;

export const coPartitionAssigner: PartitionAssigner = ({ cluster }) => ({
  name: NAME,
  version: VERSION,

  async assign({ members, topics }) {
    // implementation below
  },

  protocol({ topics }) {
    return {
      name: NAME,
      metadata: AssignerProtocol.MemberMetadata.encode({
        version: VERSION,
        topics,
        userData: Buffer.alloc(0),
      }),
    };
  },
});

protocol() adalah apa yang diumumkan setiap member waktu bergabung ke group, yang memberi tahu broker assigner mana yang didukung member itu dan topic mana yang di-subscribe. assign() menghitung siapa dapat apa, dan harus me-return memberAssignment yang di-encode lewat AssignerProtocol.MemberAssignment.encode untuk setiap member, termasuk dirinya sendiri.

Mengelompokkan berdasarkan nomor partition, bukan flat list

Round robin memperlakukan setiap pasangan topic dan partition sebagai item assignment yang terpisah. Implementasi saya justru memilih satu pemilik per index partition. Index itu dari setiap topic yang di-subscribe di-assign ke member yang sama:

typescript
async assign({ members, topics }) {
  // Deterministic member order so every member computes the same assignment.
  const sortedMembers = members.map((m) => m.memberId).sort();
  const memberCount = sortedMembers.length;

  const partitionsByTopic: Record<string, number[]> = {};
  let maxPartitions = 0;
  for (const topic of topics) {
    const ids = cluster
      .findTopicPartitionMetadata(topic)
      .map((p) => p.partitionId)
      .sort((a, b) => a - b);
    partitionsByTopic[topic] = ids;
    if (ids.length > maxPartitions) maxPartitions = ids.length;
  }

  const assignment: Record<string, Record<string, number[]>> = {};
  for (const m of sortedMembers) assignment[m] = {};

  // The crux: every topic's partition `p` is owned by the same member.
  for (let p = 0; p < maxPartitions; p++) {
    const owner = sortedMembers[p % memberCount];
    for (const topic of topics) {
      if (partitionsByTopic[topic].includes(p)) {
        (assignment[owner][topic] ??= []).push(p);
      }
    }
  }

  return sortedMembers.map((memberId) => ({
    memberId,
    memberAssignment: AssignerProtocol.MemberAssignment.encode({
      version: VERSION,
      assignment: assignment[memberId],
      userData: Buffer.alloc(0),
    }),
  }));
}

Beberapa detail yang menjelaskan implementasinya:

  • Urutan member: urutkan nilai memberId supaya daftar member yang identik menghasilkan assignment yang identik. Perubahan keanggotaan tetap bisa memindahkan banyak partition. Ini bukan sticky assignment.
  • Index partition: hitung pemiliknya dengan p % memberCount. Member itu menerima partition p dari setiap topic.
  • Input yang cocok: semua member harus subscribe ke topic relevan yang sama. Topic-topic itu harus punya jumlah partition yang sama dan key mapping yang cocok. Dengan 6 partition di market-updates dan 8 di notifications, index 6 dan 7 nggak punya pasangan. Assigner nggak bisa memperbaiki tata letak data seperti itu.

Semua ini secara manual meniru apa yang dilakukan Kafka Streams secara otomatis di dalamnya waktu Kafka Streams mensyaratkan input yang co partitioned untuk sebuah join. Tanpa library itu, kamu harus menulis perilakunya sendiri.

mermaid
graph TD
    subgraph "Round-robin (flat list, 4 pods)"
        RRList["market-p0, market-p1, ..., notif-p0, notif-p1, ..."]
        RRList --> RP1["pod-1: market-p0, market-p4, notif-p1, notif-p5"]
        RRList --> RP2["pod-2: market-p1, market-p5, notif-p2"]
        RRList --> RP3["pod-3: market-p2, notif-p0, notif-p3"]
        RRList --> RP4["pod-4: market-p3, notif-p4"]
    end
    subgraph "Co-partition assigner (by index, 4 pods)"
        CP1["pod-1 owns index 0: market-p0 AND notif-p0"]
        CP2["pod-2 owns index 1: market-p1 AND notif-p1"]
        CP3["pod-3 owns index 2: market-p2 AND notif-p2"]
        CP4["pod-4 owns index 3: market-p3 AND notif-p3"]
    end

Diagram round robin itu menggambarkan gimana nomor partition yang sama bisa sampai ke pod yang berbeda. Diagram itu bukan jejak dari satu run tertentu. Di custom assignment, member yang memiliki index p menerima partition p dari kedua topic.

Mengukur apakah assignment-nya melakukan co location

Saya menguji kedua strategi dengan 4 pod dan 6 partition per topic. Run dengan custom assigner mengonfigurasi partitionAssigners: [coPartitionAssigner] di setiap consumer. Setiap pod mencatat assignment dari GROUP_JOIN. Lalu saya menghitung apakah setiap notifikasi sampai di pod yang juga memiliki market partition-nya.

Dengan custom assigner, setiap pod punya set partition market dan notification yang cocok. Ke-30 notifikasi (10 simbol dikali 3 notifikasi masing-masing) di-join secara lokal: 30/30 lokal dan 0 remote. Run dengan round robin default menghasilkan 0 join lokal dan 30 join remote. Data dan key-nya nggak berubah di antara run.

Co partitioning aja memang nggak bakal pernah cukup. Data yang selaras di atas kertas nggak ada gunanya kalau runtime menyebarkannya ke proses yang berbeda-beda.

Memakai assignment itu untuk local join

Assignment yang selaras baru berarti kalau dipakai, jadi run kedua melakukan partition local join.

Setiap pod subscribe ke kedua topic. Pod itu membangun Map<StockSymbol, number> berisi harga terbaru dari message market-updates yang dikonsumsinya:

typescript
consumer.run({
  eachMessage: async ({ topic, message }) => {
    if (topic === MARKET) {
      const tick = tickFromKafka(message.value);
      pod.localPrice.set(tick.symbol, tick.price); // local state for owned symbols
    } else {
      const n = notifFromKafka(message.value);
      if (pod.localPrice.has(n.symbol)) pod.localJoins++;
      else pod.remoteNeeded++;
    }
  },
});

Waktu sebuah notifikasi datang, pod mengecek price map lokalnya. Di pengujian saya, custom assignment bikin pengecekan itu berhasil untuk setiap notifikasi. Tapi assignment aja nggak menjamin record market-nya udah datang lebih dulu. Kafka nggak mengurutkan record antar topic.

Setelah startup atau rebalance, aplikasi harus memulihkan atau menginisialisasi state sebelum mengandalkan local join. Dengan round robin, localPrice.has(n.symbol) juga bisa gagal karena pod lain yang memiliki market partition-nya. Pengujiannya menghitung kasus itu sebagai miss.

Map lokal bisa menghindari lookup lewat network waktu evaluasi alert. Assignment yang cocok bikin ini mungkin, tapi pemulihan state, urutan kedatangan antar topic, dan nilai yang basi tetap perlu ditangani secara eksplisit.

Padanannya di Redis, singkat aja

Redis Cluster memetakan key ke 16384 slot dengan CRC16. Hash tag di market:{BBCA} dan notif:{BBCA} memilih input hash dan slot yang sama. Tanpa tag, market:BBCA dan notif:BBCA bisa pakai slot yang berbeda.

Tag di Redis mengontrol penempatan data di node. Assigner Kafka mengontrol kepemilikan partition oleh proses consumer. Keduanya bisa mendukung operasi lokal, tapi bekerja di batasan yang berbeda.

Apa yang dijamin assigner

  • Co partitioning (jumlah partition sama, key sama, partitioner sama) menjamin tata letak datanya selaras, sedangkan consumer assignment menentukan proses mana yang membaca data itu.
  • Round robin assigner bawaan bisa menaruh partition yang cocok di consumer yang berbeda. Custom assigner menjaga keselarasannya. Kasus saya dengan 4 pod dan 6 partition termasuk salah satu yang begitu.
  • Custom PartitionAssigner butuh name, version, dan function assign yang me-return memberAssignment yang udah di-encode untuk setiap member. Cuma group leader yang terpilih yang menjalankan assign.
  • Kelompokkan partition topic berdasarkan index. Assign partition p dari setiap topic yang di-subscribe ke member yang sama.
  • Diukur di 4 pod yang sama: 30/30 notifikasi di-join secara lokal dengan assigner saya, 0/30 dengan round robin.
  • Assignment yang cocok bikin consumer bisa join notifikasi dengan price map lokalnya. Aplikasinya tetap harus menginisialisasi dan memulihkan map itu waktu kepemilikan berubah.

Semua yang sejauh ini jalan di satu broker yang nggak pernah mati. Selanjutnya: tiga broker, leader yang sengaja saya matikan, dan alasan Redis Cluster yang sangat berbeda untuk butuh tiga node.

TOPIK

SELANJUTNYA DI SERI INIKafka ISR vs Redis Cluster: Replication dan Failover

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.