"Exactly-once" di Kafka punya batas transaction yang spesifik. Kafka bisa meng-commit output record dan consumer offset sebagai satu operasi atomic. Jaminan ini nggak mencakup write ke database eksternal, panggilan pembayaran, atau email.
Saya menyelidiki ini waktu menambahkan price alert ke ticker saya. Sebuah tick melewati threshold, consumer mem-produce sebuah notification record, dan source offset-nya maju. Saya bikin versi transactional untuk menguji apakah record dan offset itu bisa di-commit bersamaan.
Apa arti exactly once
Delivery semantics pada dasarnya soal kapan kamu commit dibandingkan kapan kamu mengerjakan pekerjaannya:
| Semantic | Cara | Failure mode |
|---|---|---|
| At-most-once | Commit sebelum memproses | Crash setelah commit, sebelum pekerjaan selesai: message hilang |
| At-least-once | Commit setelah memproses | Crash setelah pekerjaan selesai, sebelum commit: message diproses ulang |
| Exactly-once | Output dan commit offset dalam satu transaction | Nggak satu pun terjadi, atau keduanya terjadi: nggak ada efek parsial |
Dengan commit setelah memproses, crash bisa menyebabkan handling yang berulang. Karena itu, efek yang idempotent jadi wajib. Exactly once semantics (EOS) di Kafka berlaku untuk pipeline yang lebih sempit: consume Kafka record, produce Kafka record, dan commit source offset dalam satu transaction.
Kalau transaction itu di-abort, committed reader nggak melihat output-nya maupun offset-nya yang maju. Jaminan ini bergantung pada identitas producer, penanganan transaction, dan isolation consumer yang benar.
EOS mencakup Kafka ke Kafka. Kalau langkah "process" kamu memanggil API pembayaran, menulis ke Postgres, atau mengirim email, nggak satu pun dari itu ada di dalam transaction. Transaction-nya bisa di-abort setelah side effect kamu udah terjadi, atau di-commit setelah side effect itu gagal. Kafka nggak punya cara untuk me-rollback sebuah HTTP call. EOS cocok untuk pipeline yang consume dari Kafka dan produce ke Kafka, misalnya mengubah tick market-updates jadi notifications waktu sebuah price alert terpicu. EOS nggak bisa bikin write ke dua sistem eksternal jadi atomic.
Idempotent producer: PID dan sequence number
Idempotent producer menangani produce request yang berulang. Sebuah request bisa aja berhasil di broker, tapi acknowledgment-nya baru sampai setelah client timeout. Client lalu me-retry batch yang sama. Idempotence di producer mencegah retry itu menambahkan batch yang sama dua kali.
export function createIdempotentProducer(): Producer {
return kafka.producer({
idempotent: true, // forces acks=all and dedup on retry
maxInFlightRequests: 5, // max allowed while keeping idempotence/ordering
retry: { retries: 10 },
});
}
Dengan idempotent: true, KafkaJS mewajibkan acks=all dan producer-nya dapat sebuah Producer ID (PID). Setiap batch juga membawa sequence number untuk partition-nya. Broker memakai identitas dan informasi sequence ini untuk mengenali pengiriman yang berulang. Broker bisa menolak retry yang kalau dibiarkan bakal menambahkan batch yang sama sekali lagi.
Idempotence di producer butuh setting acknowledgment, retry, dan request concurrency yang kompatibel. Jangan berasumsi bahwa mengaktifkannya otomatis men-set maxInFlightRequests ke 5 di KafkaJS. Panduan transaction KafkaJS menetapkan maxInFlightRequests: 1, idempotent: true, dan sebuah transactionalId untuk konfigurasi EOS-nya.
Idempotence mencakup retry dari request satu producer. Kode aplikasi tetap bisa mengeluarkan logical event yang duplikat, dan jaminan ini nggak mencakup beberapa partition maupun beberapa topic sebagai satu kesatuan. Pemrosesan di consumer juga bisa berulang setelah crash sebelum offset di-commit. Idempotent produce dan exactly once consume produce adalah jaminan yang berbeda, dan itulah kenapa transaction ada sebagai layer terpisah yang opt in di atasnya.
Transaction: baca, proses, tulis, secara atomic
Transactional producer butuh transactionalId yang stabil untuk logical input assignment-nya. Dengan begitu, Kafka bisa menolak instance producer yang lebih lama setelah penggantinya mulai jalan. Producer aktif untuk assignment yang berbeda nggak boleh berbagi identitas yang sama.
const producer = kafka.producer({
transactionalId: "s04-eos",
idempotent: true,
maxInFlightRequests: 1,
});
const txn = await producer.transaction();
try {
await txn.send({ topic: "notifications", messages });
// The defining feature of EOS: move the input group's offsets *inside*
// the same transaction as the output. Crash here and neither side commits.
await txn.sendOffsets({
consumerGroupId: "s04-eos-consumer",
topics: [{ topic: "market-updates", partitions: [/* ... */] }],
});
await txn.commit();
} catch (e) {
await txn.abort();
throw e;
}
Panggilan sendOffsets memasukkan consumer offset yang baru ke dalam transaction output. Kafka menyimpan group offset di __consumer_offsets. Transaction-nya meng-commit output dan offset bersamaan, atau nggak meng-commit keduanya. Ini menghilangkan celah antara output Kafka yang berhasil dan majunya source offset yang bersangkutan.
Commit vs abort: melihat sebuah batch menghilang
Untuk menguji visibility transaction, saya membuat dua batch yang masing-masing berisi lima notifikasi. Setiap batch mencakup persilangan threshold harga untuk lima simbol yang berbeda:
Transaction A mengirim lima notifikasinya, memanggil sendOffsets untuk memajukan offset dari source consumer group, lalu commit.
Transaction B mengirim lima notifikasi, lalu abort.
sequenceDiagram
participant P as Producer
participant TC as Coordinator
participant N as notifications topic
participant C as read_committed Consumer
P->>TC: begin transaction
P->>N: send 5 notifications
P->>TC: sendOffsets for consumer group
P->>TC: commit
TC->>N: write commit marker
C->>N: poll
N-->>C: 5 notifications visible
Lalu saya membaca notifications dengan consumer read_committed dan read_uncommitted. Default KafkaJS adalah readUncommitted: false, tapi default-nya beda-beda di setiap client. Set isolation mode secara eksplisit kalau kamu pakai client lain. Hasilnya:
read_committed: 5 notifikasi. Cuma batch dari transaction A.read_uncommitted: 10 notifikasi. Kedua batch, termasuk yang di-abort.
graph LR
TB["Transaction B: 5 notifications"] --> Abort
Abort --> RC["read_committed: invisible"]
Abort --> RU["read_uncommitted: visible"]
Kafka menulis record yang di-abort ke log, diikuti transaction marker. Consumer dengan read_committed mengecualikan record yang di-abort, sedangkan consumer dengan read_uncommitted bisa mengembalikannya. Hasil 5 lawan 10 itu menunjukkan perbedaan visibility ini. Hasil itu nggak berarti Kafka menghapus byte yang di-abort dari storage.
Transaction menambah request ke coordinator. Producer di atas juga memakai maxInFlightRequests: 1, yang membatasi request yang berjalan bersamaan. Pakai desain ini kalau output Kafka dan source offset harus di-commit bersamaan. Efek eksternal tetap butuh idempotency atau koordinasi transaction tersendiri.
Atomicity command Redis punya batasan yang berbeda
Redis bisa mengeksekusi XADD untuk output dan XACK untuk input dalam satu Lua script atau transaction MULTI/EXEC. Di Redis Cluster, key yang terkait harus berada di slot yang sama. Hash tag seperti market:{BBCA} bisa memberikan penempatan itu.
Ini adalah eksekusi command Redis yang atomic, dengan aturan error dan recovery yang berbeda dari Kafka transaction. Aplikasinya harus menangani deduplication, pemrosesan, dan durability secara eksplisit. Transaction Redis nggak me-rollback command setelah terjadi execution error. Lihat dokumentasi transaction Redis.
Kalau pemrosesan dan acknowledgment terpisah, crash sebelum XACK bisa menyebabkan pekerjaan berulang. Consumer lain mungkin mengambil alih entry itu dengan XAUTOCLAIM. Handler-nya harus tahan terhadap pengulangan itu. Redis maupun Kafka sama-sama butuh koordinasi tambahan kalau pemrosesannya mengubah sistem eksternal.
Apa artinya buat kamu
Aktifkan idempotence di producer dengan setting yang dibutuhkan client-nya. Pakai Kafka transaction kalau output record dan source offset butuh satu batas commit. Ukur biaya coordinator dan concurrency untuk workload-nya. Untuk write ke database dan panggilan API eksternal, rancang idempotency atau koordinasi eksplisit di batas sistem itu.
Post berikutnya membandingkan Redis Pub/Sub dengan Streams, termasuk apa yang terjadi waktu sebuah subscriber terputus.

Memuat komentar...