- Dalam verifikasi Bufstream 0.1.0 hingga 0.1.3, sistem streaming yang kompatibel dengan Kafka, ditemukan 2 masalah ketersediaan dan 3 masalah keamanan pada Bufstream sendiri, dan per 0.1.3 kelima masalah tersebut telah diperbaiki
- Pengujian didasarkan pada Java Kafka Client 3.8.0 serta pengujian Jepsen Kafka/Redpanda yang sudah ada, dengan menggunakan konfigurasi yang mengutamakan keamanan seperti
acks = all,enable.idempotence = true,enable.auto.commit = false, danread_committed - Masalah Bufstream mencakup penghentian konsumen dan produsen, respons offset
0yang keliru, hilangnya commit transaksi, serta hilangnya penulisan yang telah diakui akibat bug pemfilteran ukuran respons fetch API - Dalam proses investigasi, pada Kafka Java client dan protokol transaksi Kafka juga terungkap masalah seperti
Consumer.close()yang memblokir tanpa batas, offset konsumen yang tidak dapat diprediksi, serta aborted read, lost write, dan torn transaction - Jepsen menilai bahwa protokol transaksi Kafka tidak secara eksplisit menjamin urutan permintaan klien dan nomor transaksi, sehingga saat menggunakan Java client resmi, keamanan transaksi Kafka dan sistem yang kompatibel dengan Kafka dapat rusak
Struktur Bufstream dan cakupan verifikasi
- Kafka adalah sistem streaming yang menyediakan log append-only yang direplikasi dan di-shard, sementara Bufstream adalah implementasi alternatif Kafka yang memprioritaskan data governance dan efisiensi biaya di lingkungan cloud
- Bufstream menyediakan topic dan partition seperti Kafka, serta bekerja dengan klien Kafka standar
- producer menambahkan record dengan
producer.send() - consumer diikat ke partition dengan
consumer.assign()atauconsumer.subscribe(), lalu membaca record denganconsumer.poll() - consumer group membagi pemrosesan record dari sekumpulan topic
- producer menambahkan record dengan
- Jika diintegrasikan dengan Buf Schema Registry, Bufstream dapat memeriksa record Protocol Buffer untuk mendukung validasi record, field-level access control, dan konversi format data dengan sistem lain
- Berbeda dengan Kafka yang menggunakan disk lokal dan protokol replikasi sendiri, Bufstream menulis data langsung ke object storage
- Tujuannya adalah menekan biaya dengan memanfaatkan struktur biaya trafik replikasi object storage
- Node Bufstream dapat berjalan sebagai VM stateless yang di-auto-scale
- Bufstream terdiri dari tiga subsistem
- agent: layanan stateless yang menyediakan Kafka API
- object store: menyimpan chunk record dan menyediakannya untuk reader
- coordination service: saat ini menggunakan etcd, menentukan chunk mana yang telah di-commit dan urutan record
- Per Oktober 2024, Bufstream hanya diterapkan untuk sebagian pelanggan; dokumentasinya menonjolkan “drop-in replacement untuk Apache Kafka” serta kompatibilitas dengan Kafka transactions dan exactly-once semantics, tetapi tidak banyak memuat klaim keamanan yang spesifik
Konfigurasi klien dan asumsi transaksi
- Seperti pengujian sebelumnya terhadap sistem yang kompatibel dengan Kafka, Jepsen menyesuaikan konfigurasi client untuk mendapatkan perilaku yang lebih aman
-
Konfigurasi Producer
- Menggunakan nilai default
acks = all - Di Bufstream,
acks = 0dapat mengakui penulisan tanpa menunggu storage, sehingga penulisan yang sudah di-commit bisa hilang acks = 1danacks = allmemblokir hingga Bufstream yakin data telah dipersistenkan secara durable- Untuk mencegah append duplikat akibat retry otomatis Kafka producer, digunakan nilai default
enable.idempotence = true
- Menggunakan nilai default
-
Konfigurasi Consumer
- Karena ada dokumentasi yang menyebut auto-commit dapat menyebabkan kehilangan data, umumnya digunakan
enable.auto.commit = false - Jika tidak ada committed offset, nilai default
auto.offset.resetdimulai dari offset terbaru, sehingga tidak menjamin at-least-once delivery - Agar consumer dapat mengamati seluruh log, digunakan
auto.offset.reset = earliest - Transaksi Kafka terdiri dari kumpulan record yang dikirim producer dan map offset maksimum per partition yang di-poll oleh consumer
- Hanya ketika transaction di-commit, record yang dikirim menjadi durable dan pada akhirnya terlihat oleh consumer
read_committed, sementara committed offset juga naik ke setidaknya offset yang ditentukan dalam transaction - Jika transaction tidak di-commit, committed offset tidak maju, dan visibilitas penulisan dapat berbeda tergantung konfigurasi consumer
- Fenomena consumer
read_uncommittedmembaca nilai dari transaction yang diabort diklasifikasikan sebagai aborted read (G1a) - Dokumentasi Kafka menyatakan bahwa
read_committedmencegah G1a dan sampai batas tertentu menjamin bahwa semua penulisan dalam transaction terlihat atau tidak ada yang terlihat sama sekali, tetapi dalam pengujian Jepsen terhadap Kafka, Redpanda, dan Bufstream, teramati write cycle (fenomena mirip G0) dan beberapa bentuk G1c
- Karena ada dokumentasi yang menyebut auto-commit dapat menyebabkan kehilangan data, umumnya digunakan
Desain pengujian
- Jepsen menguji Bufstream 0.1.0 hingga 0.1.3 serta beberapa build release candidate
- Test harness menggunakan Bufstream test harness, Jepsen testing library, dan Java Kafka Client 3.8.0
-
Lingkungan eksekusi
- Menggunakan 3 hingga 5 node Debian Bookworm, baik sebagai container LXC maupun VM EC2
- 1 node untuk etcd, 1 node untuk Minio, dan sisanya sebagai agent Bufstream
- Producer, consumer, dan admin client diinisialisasi dengan hanya satu node di
bootstrap_servers, tetapi smart client discovery tidak diblokir
-
Konfigurasi keamanan utama
- auto-commit false
acks = all- retries 1.000
- idempotence enabled
- isolation level
read_committed auto_offset_reset = earliest- pembuatan topic otomatis di sisi server dinonaktifkan
- Injeksi gangguan mencakup process pause (
SIGSTOP), crash (SIGKILL), clock skew (clock_settime), dan network partition (iptables) - Karena Bufstream terbagi menjadi agent, object store, dan coordination service, Jepsen membuat alat baru yang dapat menginjeksikan gangguan hanya ke subsistem tertentu
- Misalnya, kombinasinya diubah seiring waktu, seperti hanya melakukan crash pada node Bufstream atau hanya melakukan pause pada coordinator etcd
Queue workload dan Abort workload
- Queue workload menganalisis keamanan sesuai dengan model data Kafka
- Setiap logical process menjalankan producer, consumer, dan admin client
- numeric key mengidentifikasi topic-partition tertentu
- key dipilih dengan exponential frequency sehingga sebagian key sering diakses, sementara sebagian lainnya jarang
- Menggunakan tiga operation dasar
crash: menghentikan logical process dan menggantinya dengan client barusubscribeatauassign: mengubah kumpulan topic atau partition yang akan di-polloleh consumertxn,poll,send: menjalankan sequence micro-operationpollatausend
- Pada non-transactional workload, setiap
sendataupollhanya berisi tepat satu micro-operation - Pada transactional workload, beberapa micro-operation dibungkus sebagai Kafka transaction
- Analisis membuat mapping offset-to-value per key, lalu mencari error
- Jika beberapa value terlihat pada offset yang sama, disebut inconsistent offset
- Jika value yang sama terlihat pada beberapa offset, disebut duplicate error
- Jika record yang diakui sama sekali tidak teramati, dikategorikan sebagai lost atau unseen
- Jika poll mengembalikan value yang dikirim oleh operation yang diabort, disebut aborted read
- Juga memeriksa apakah transaction mengamati tulisannya sendiri
- Setelah main test, kegagalan diselesaikan dan masuk ke tahap final reads
- Setiap process membaca semua topic-partition mulai dari offset 0 dan melakukan poll hingga written offset tertinggi yang diketahui
- Jika final reads timeout dan record yang diakui masih belum teramati, record diklasifikasikan sebagai unseen
- Abort workload ditambahkan untuk melacak perilaku poll offset setelah transaction abort
- topic dibatasi pada satu partition, process, producer, dan consumer
- Setelah transaction melakukan poll record lalu sengaja diabort, poll offset berikutnya diklasifikasikan sebagai advance, rewind, rewind-further, atau other
5 masalah yang ditemukan di Bufstream
-
Consumer macet (#1)
- Dari 0.1.0 hingga 0.1.3-rc.8, tahap final read sering macet
consumer.poll()langsung mengembalikan hasil kosong, tetapi di log masih tersisa ribuan record yang sudah diakui- Kondisi ini berlangsung dari puluhan detik hingga lebih dari 1 jam
- Dalam satu pengujian, selama 120 detik pertama dikirim 691 record yang diakui, dan pada saat final reads dimulai, 40 di antaranya tidak teramati oleh poller mana pun
- Setelah itu, selama lebih dari 1 jam
consumer.poll()tidak mengembalikan hasil sehingga pengujian timeout - Penyebabnya adalah node Bufstream yang restart dapat mengembalikan nilai cache stale untuk last stable offset dan high watermark
- Sebagian client library menyimpulkan tidak ada record lebih lanjut lalu stall, dan Bufstream menerapkan patch pada 0.1.3-rc.6 agar cache di-refresh saat startup
-
Producer & consumer macet (#2)
- Bahkan pada 0.1.3-rc.6, masalah unseen write setelah pause, crash, dan partition pada coordinator, storage, serta node Bufstream masih terus teramati
- Dalam beberapa kasus, setelah coordinator pause, meski semua node Bufstream sedang berjalan, client masuk ke kondisi menunggu
InitProducerIdhingga timeout - Dalam kasus lain,
listOffsetsgagal dengannode ... being disconnectedatautimed out waiting for a node assignment, sementarapollselesai tetapi tidak mengembalikan hasil - Jika node Bufstream di-kill lalu di-restart, masalahnya hilang
- Penyebabnya terkait etcd lease
- Agen Bufstream menggunakan etcd leases untuk melacak agen aktif
- Karena pause atau partition singkat, etcd menghapus key yang terikat pada lease agen, tetapi update penghapusan itu bisa tidak terkirim ke agen
- Agen menjadi tidak tahu bahwa lease miliknya telah hilang
- Tim Bufstream menambahkan logic polling tambahan, dan pada 0.1.3-rc.8 unseen write sebagian besar teratasi
-
Offset nol palsu (#3)
- Dari 0.1.0 hingga 0.1.3-rc.2, value yang dikirim bisa diberi offset
0, lalu muncul pada offset sebenarnya yang lebih tinggi - Ini terjadi bahkan ketika offset
0sudah dialokasikan jauh sebelumnya - Hanya sender yang mengamati offset 0, sedangkan poller mengamati offset yang lebih tinggi
- Dalam pengujian 2 menit dengan satu node Bufstream dan pause pada proses etcd, 6 write mendapat offset
0lalu muncul pada offset yang lebih tinggi - Penyebabnya adalah field yang diperlukan tidak ada dalam error response Bufstream
- Bufstream mengirim permintaan log commit ke etcd dan etcd memprosesnya, tetapi karena pause atau partition, Bufstream bisa timeout saat menunggu response
- Bufstream mengirim error code ke client, tetapi tidak menetapkan offset record yang dikirim ke
-1, yaitu sinyal error - Java Kafka client menafsirkannya sebagai response sukses untuk offset
0 - Franz-go, yang digunakan test suite Bufstream, menafsirkan message ini sebagai error sehingga masalah tersebut tidak muncul dalam pengujian
- Bufstream memperbaikinya pada 0.1.3-rc.6, dan setelah itu Jepsen tidak dapat mengamatinya lagi
- Dari 0.1.0 hingga 0.1.3-rc.2, value yang dikirim bisa diberi offset
-
Write transaction hilang (#4)
- Pada 0.1.2, sebagian record dari transaction yang sudah commit menghilang dan tidak pernah teramati lagi; write loss ini sering terjadi
- Dalam satu pengujian, selama 100 detik dan 6.761 write transaction, 240 record yang ditulis oleh transaction yang commit hilang
- Dalam contoh, value
141untuk key5dikembalikan sebagai berhasil ditulis pada offset274, tetapi semuaconsumer.poll()melewati offset tersebut - Penyebabnya adalah bug pada concurrency safety mechanism yang ditambahkan di 0.1.2
- Mechanism ini memberi nomor unik pada setiap transaction dalam producer epoch untuk mengurangi kurangnya idempotence pada Kafka transaction protocol
- Karena bug pada logic pelacakan transaction number, ketika beberapa transaction di berbagai epoch di-commit, sebagian commit keliru diabaikan
- Transaction yang tampak sudah commit bisa sebenarnya di-abort, atau sebaliknya
- Jepsen menemukan bug ini berkat transaction timeout yang disetel rendah, yaitu 1 detik
- Bufstream mengidentifikasi masalah ini dalam beberapa jam setelah rilis 0.1.2, mencegah customer melakukan upgrade, dan customer tidak melakukan upgrade ke 0.1.2
- Perbaikannya disertakan dalam 0.1.3-rc2
-
Lost writes akibat filtering di sisi server (#5)
- Pada 0.1.3-rc.8, setelah gangguan kecil seperti pause pada process Bufstream atau coordinator, atau partition di antara keduanya, window write loss singkat sering muncul
- Data loss terjadi terlepas dari apakah transaction digunakan atau tidak
- Dalam satu pengujian 5 menit, dari 16.770 record, 22 record di-acknowledge tetapi tidak dapat di-poll oleh consumer mana pun
- Sebagian record sempat terlihat oleh poller untuk beberapa waktu, tetapi kemudian menghilang dari poll
- Penyebabnya adalah logic pembatasan ukuran response fetch API yang ditambahkan pada 0.1.3-rc.8 untuk mengakali bug pada Kafka web GUI yang populer
- Bug pada filtering logic menyembunyikan record dari lagging consumer, sehingga tampak seperti write loss
- Bufstream memperbaikinya pada 0.1.3-rc.12
Masalah pada Kafka Java client dan protokol Kafka
-
KIP-588: ProducerFencedException yang menyesatkan
- Selama pengujian, error
ProducerFencedException: There is a newer producer with the same transactionalId which fences the current one.sering terjadi - Error ini muncul bahkan dalam pengujian ketika semua producer mendapat transactional ID yang unik, sehingga butuh waktu untuk menemukan penyebabnya
- KIP-588 menyebutkan bahwa
ProducerFencedExceptionjuga dapat dilempar saat transaction timeout - Kafka Java client menggunakan
TimeoutExceptionkhusus untuk sebagian besar timeout, tetapi dalam kasus ini melemparProducerFencedException - Padahal sebenarnya tidak ada producer yang bertabrakan, tetapi pesan error mengatakan ada instance producer kedua
- KIP-588 sudah terbuka selama 2 tahun, dan Jepsen merekomendasikan tim Kafka mengubah error message tersebut
- Selama pengujian, error
-
KAFKA-17734: Consumer.close() dapat block tanpa batas
- Dalam pengujian di Bufstream maupun Kafka, tes berhenti setiap beberapa jam karena bug Java client
Consumer.close()secara default melakukan block pada network IO- Parameter timeout pada
close()seharusnya mencegah block tanpa batas, tetapi tidak berfungsi - Cara memanggil
consumer.wakeup()dari thread terpisah untuk menginterupsi consumer yang stuck di IO juga tidak efektif - Jepsen menilai program jangka panjang harus dapat melepaskan resource seperti client, connection, thread, dan memory dalam waktu yang wajar meski terjadi network error, lalu mendaftarkan KAFKA-17734
-
KAFKA-17582: offset consumer tidak dapat diprediksi setelah transaction gagal
- Dokumentasi resmi Kafka hampir tidak menjelaskan bagaimana consumer offset seharusnya saat transaction commit gagal
- Dokumentasi desain Kafka dari Confluent mengatakan bahwa jika transaction di-abort, consumer position kembali ke nilai sebelumnya, tetapi Java client sebenarnya tidak selalu berperilaku demikian
- Dari hasil workload abort, bahkan pada cluster yang sehat, perilaku setelah abort sulit diprediksi
- Sebagian besar transaction pair maju ke offset yang lebih jauh
- Sebagian lainnya rewind ke offset sebelumnya
- Semua rewind terkait dengan rebalance event, dan semua advance tidak disertai rebalance
- Menurut jawaban dari pihak Kafka, perilaku ini intentional
- Consumer terus advance
- Jika rebalance terjadi, consumer dapat rewind ke titik sembarang sesuai committed offset
- Pengguna harus me-rewind consumer position secara manual saat transaction abort
- Jepsen membuka KAFKA-17582, serta mengusulkan dokumentasi perilaku ini dan peninjauan perubahan default rewind saat transaction abort
- Queue workload juga diperbaiki agar secara eksplisit me-rewind consumer
-
KAFKA-17754: write loss, aborted read, torn transaction
- Pada Bufstream 0.1.0~0.1.3, hanya dengan pause proses Bufstream, pause coordinator, crash, dan network partition, ditemukan aborted read, lost write, dan pelanggaran atomicity
- Analisis mengarah pada cacat mendasar dalam Kafka transaction protocol
- Dalam contoh, client menjalankan transaction dengan transactional ID unik
jt1234dan mengirimcommitted = falsekeEndTxnuntuk abort, tetapi 15 pemanggilanpoll()mengamati write dari transaction yang sudah di-abort - Write lain dari transaction yang sama tidak diamati oleh poller mana pun
- Jika packet capture dan log Bufstream dilihat bersama, penyebabnya adalah commit message yang tertunda
- Commit
EndTxnyang dikirim beberapa transaction sebelumnya diproses terlambat di satu node - Client sudah melanjutkan transaction berikutnya
- Commit yang tertunda diterapkan pada transaction saat ini, sehingga hanya bagian awal transaction yang ter-commit, sementara sisanya diperlakukan seperti transaction terpisah dan di-abort
- Kafka protocol dirancang agar client dapat mengirim request ke beberapa TCP connection dan beberapa node, tetapi tidak memiliki sequence number untuk menentukan urutan request dari client yang sama
- Tidak ada pula konsep transaction number, sehingga ketika server menerima commit atau abort message, server tidak tahu transaction mana yang hendak diakhiri oleh client
- Akibatnya, situasi berikut menjadi mungkin
- Transaction yang tampak committed sebenarnya di-abort
- Transaction yang di-abort sebenarnya di-commit
- Terjadi torn transaction, yaitu hanya sebagian write dalam transaction yang dipertahankan dan sebagian lainnya hilang
- Kafka Java client resmi memperlakukan timeout sebagai retryable dan dapat otomatis mengirim beberapa message
EndTxn, sehingga masalah dapat terjadi meski pengguna hanya memanggil commit atau abort satu kali per transaction - Jepsen juga mengamati aborted read dan torn transaction pada Kafka akibat process pause, lalu membuka KAFKA-17754
- Para engineer Kafka menilai KIP-890 kemungkinan dapat memperbaiki masalah ini
- KIP-890 mengubah transaction protocol dengan menaikkan producer epoch pada setiap transaction
- Karena server menolak message dari epoch sebelumnya, commit message dari transaction lama dapat dicegah agar tidak bocor masuk ke transaction berikutnya
- Pada 0.1.3, Bufstream menambahkan mechanism yang menggunakan etcd revision sebagai logical clock untuk mengurangi frekuensinya, tetapi belum mencegah reorder antara client dan Bufstream
- Jepsen tetap mengamati aborted read, lost write, dan torn transaction pada 0.1.3, dan menilai diperlukan penyelesaian di sisi client
Ringkasan Hasil Keseluruhan
- Kelima masalah pada Bufstream sendiri semuanya telah diperbaiki
- #1: consumer stuck karena lagging highest stable offset, tidak perlu gangguan, diperbaiki di 0.1.3-rc.6
- #2: producer/consumer stuck karena kedaluwarsa lease etcd, perlu pause, diperbaiki di 0.1.3-rc.8
- #3: offset nol palsu, perlu pause, diperbaiki di 0.1.3-rc.6
- #4: penulisan transaksi hilang, tidak perlu gangguan, diperbaiki di 0.1.3-rc.2
- #5: penulisan hilang karena filtering sisi server, perlu pause, diperbaiki di 0.1.3-rc.12
- Masalah terkait Kafka masih tetap ada
- KIP-588: pesan error yang keliru pada transaction timeout, belum terselesaikan
- KAFKA-17734:
ConsumerClient.close()dapat block tanpa batas waktu, belum terselesaikan - KAFKA-17582: offset consumer tidak dapat diprediksi setelah kegagalan transaksi, belum terselesaikan
- KAFKA-17754: write loss, aborted read, torn transaction, belum terselesaikan
- Jepsen mengingatkan bahwa verifikasi keamanan eksperimental dapat membuktikan keberadaan bug, tetapi tidak dapat membuktikan ketiadaannya
- Secara khusus, Jepsen menilai bahwa karena KAFKA-17754, sulit menentukan apakah ada kasus write loss lain di Bufstream
Rekomendasi untuk Pengguna dan Operasional Bufstream
- Pengguna yang memakai transaksi Bufstream dengan Java Kafka client resmi perlu mempertimbangkan bahwa transaksi saat ini mungkin tidak aman
- Transaksi yang di-abort bisa saja benar-benar ter-commit
- Transaksi yang di-commit bisa saja benar-benar ter-abort
- Transaksi dapat terbelah, sehingga hanya sebagian efeknya yang dipertahankan
- Bufstream memandang client Franz-go kurang rentan terhadap masalah ini, tetapi Jepsen tidak menguji Franz-go dengan teknik seperti pada pekerjaan ini
- Client lain bisa saja rentan ataupun tidak
- Pengguna sebelum Bufstream 0.1.3 dapat mengalami masalah berikut
producer.send()keliru mengembalikan offset0, bukan offset sebenarnya- masalah availability metastable yang membuat client stuck
- Jepsen merekomendasikan upgrade ke 0.1.3
- Jepsen menilai arsitektur keseluruhan Bufstream tampak sound
- Cara menentukan urutan chunk data immutable menggunakan coordination service seperti etcd adalah pendekatan yang relatif sederhana dan memiliki preseden di OLTP serta streaming system
- Dari sisi operasional, ada dua perbaikan yang direkomendasikan
- Jika permintaan shared file ke storage gagal saat startup, cluster dapat crash; Jepsen merekomendasikan penambahan retry, dan Bufstream telah menambahkan retry layer
- Saat dependency unavailable, Jepsen merekomendasikan agar agent tidak langsung mati, melainkan tetap berjalan, menyediakan backpressure dan system status, serta recover dengan lebih mulus
- Per 0.1.3, Bufstream telah menambahkan retry logic tambahan untuk etcd, tetapi masih memerlukan constant supervision agar tetap online
- Pengguna harus memastikan ada process supervisor dan menguji apakah proses tetap berjalan tanpa menyerah selama outage jangka panjang
Perlunya Dokumentasi dan Perbaikan Protokol Transaksi Kafka
- Dokumentasi resmi Kafka hampir tidak membahas transaksi, sehingga pengguna harus menggabungkan berbagai source yang ambigu dan saling bertentangan
- Jepsen merekomendasikan tim Kafka membuat dokumen terpusat yang merangkum semantics transaksi secara jelas, dan menyebut KAFKA-17671
- Dokumen tersebut setidaknya harus menjelaskan hal-hal berikut
- kapan consumer mengamati offset yang meningkat secara monotonik
- kapan consumer dapat melewati record yang telah di-acknowledge
- apakah rebalance dapat memengaruhi transaksi di tengah proses
- kapan offset write producer meningkat secara monotonik
- kapan G0, G1a, G1b, G1c, fractured read, dan pembacaan write dari transaksi sendiri dianggap legal
- apa arti nilai kembalian
poll()dan offset setelah transaksi yang di-abort - bagaimana menangani transaction error, error saat abort, dan error saat rewind
- Dokumentasi Confluent berulang kali mengatakan bahwa default Kafka menyediakan at-least-once delivery, tetapi Jepsen menunjukkan bahwa ini tampaknya tidak benar
auto.offset.reset = latestdapat membuat record yang belum diproses seolah-olah sudah “committed”- Dokumentasi offset management Confluent juga menyebut risiko kehilangan message progress saat crash pada auto-commit default
- Dokumentasi yang menyatakan consumer di-rewind saat transaction abort juga berbeda dari kenyataannya
- Jepsen menilai protokol transaksi Kafka perlu diperbaiki secara mendasar
- Protokol secara implisit mengasumsikan ordered reliable delivery, tetapi ada process pause, network unreliability, non-zero latency, dan unordered delivery di antara beberapa TCP socket
- Protokol Kafka mendistribusikan message ke beberapa node dan TCP socket, sementara client otomatis melakukan retry message
- Tidak ada sequence number untuk memulihkan urutan message dari client yang sama, maupun transaction number untuk memverifikasi target transaksi
- KIP-890 mencoba menjamin urutan yang lebih ketat dengan menaikkan epoch pada setiap transaction commit
- Client library juga dapat membantu dengan me-re-initialize producer untuk menaikkan epoch ketika message tidak di-acknowledge
- Java Kafka Client 3.8.0 rentan terhadap masalah ini
- Jepsen menilai Franz-go dapat memitigasi atau mencegah masalah dengan melakukan re-initialize saat timeout, tetapi belum meneliti client library lain
Pekerjaan Mendatang
- Banyak pengguna lebih mengandalkan “exactly-once semantics” dari Kafka Streams API daripada menangani transaction secara langsung, sehingga di masa mendatang akurasi aplikasi Streams dapat diteliti
- Jepsen juga menemukan unseen write di Kafka saat menyelidiki KAFKA-17754, tetapi tidak sempat menganalisisnya karena keterbatasan waktu
- unseen write bisa menjadi tanda hanging transaction, stuck consumer, atau data loss
- Masih menjadi pertanyaan apakah message
Produceyang tertunda dapat masuk ke transaction di masa depan dan melanggar transaction guarantee - Jepsen juga mencurigai kemungkinan Kafka Java Client menggunakan ulang sequence number saat request timeout, sehingga write di-acknowledge tetapi diam-diam di-discard
- Ketika rebalance event terjadi, consumer position dapat bergerak maju atau mundur, tetapi aturannya tidak jelas
- Jika Kafka mendokumentasikan perilaku yang dimaksudkan, Jepsen ingin memverifikasinya
- Jepsen menjelaskan bahwa karena merupakan random process, sulit mengeksplorasi anomaly yang jarang terjadi
- Masalah yang hanya terjadi sekali sangat sulit untuk debugging dan reproduction
- Bufstream juga menggunakan Antithesis, yang menjalankan seluruh distributed system dalam deterministic hypervisor dan simulated network
- Menggabungkan workload generation dan history checking dari Jepsen dengan deterministic, replayable environment milik Antithesis dapat meningkatkan reproducibility pengujian
1 komentar
Komentar Hacker News
Saat menyelidiki isu seperti KAFKA-17754, kalau Jepsen juga menemukan penulisan yang tidak terlihat di Kafka, sepertinya sudah waktunya Jepsen menguliti Kafka lagi secara mendalam
Penyelidikan terakhir dilakukan pada 2013 (https://aphyr.com/posts/293-call-me-maybe-kafka, Kafka 0.8 beta), dan sekarang tampaknya Kafka sendiri baru berada pada tahap mulai menemukan berbagai masalah
Hal seperti “penulisan sudah dikonfirmasi tetapi bisa dibuang diam-diam” cukup menakutkan
Sangat mengejutkan bahwa pada nilai bawaan
enable.auto.commit=true, konsumen Kafka bisa meng-commit offset terlepas dari apakah aplikasi benar-benar telah memprosesnyaSaya tidak pernah memahami auto-commit seperti itu, dan menurut saya nilai bawaan seperti itu tidak masuk akal
Penjelasan dokumentasinya memang tidak terlalu jelas, tetapi secara keseluruhan terbaca seolah offset hanya di-commit setelah pemrosesan selesai
Saya memahaminya bahwa pengaturan interval auto-commit, seperti yang diharapkan dalam at-least-once processing, membantu mengurangi jendela pemrosesan duplikat, bukan kehilangan pesan
Jika tidak melakukan commit secara eksplisit, Kafka tidak punya cara untuk tahu apakah pesan telah diproses
Kafka mengasumsikan pesan yang diserahkan langsung diproses
Auto-commit mirip seperti menyerahkan es krim cone lalu langsung berbalik dan menganggap orang itu sudah memakannya. Bisa saja ada orang yang menjatuhkannya begitu menerima dan belum sempat menggigit sedikit pun
Jika menginginkan jaminan itu, Anda harus melakukan acknowledgement secara eksplisit
Misalnya, jika Anda hanya menulis pesan ke database, pesan itu dianggap terkonfirmasi saat masuk ke callback handler klien
Namun dalam praktiknya, Anda kemungkinan besar ingin pengonfirmasian baru terjadi setelah insert DB berhasil
DB bisa menjadi tidak dapat diakses karena jaringan, Kubernetes, pengaturan firewall, dan sebagainya, lalu di tengah situasi itu engineer mencoba me-restart sehingga klien mati, dan pesan yang belum diproses pun mudah muncul
Sistem lain bisa menentukan apakah terjadi kegagalan, dan fitur ini dapat menggeser batas atas untuk mengurangi pemrosesan ulang
Namun jika timing-nya pas dan terjadi gangguan, Anda harus berasumsi bahwa setelah restart Anda bisa menerima ulang sebagian data yang sebenarnya sudah diproses
Masalahnya adalah ketika tidak ada pemrosesan seperti itu sebelum auto-commit
Jika dibaca, kesannya memang dirancang agar commit terjadi cukup lama setelah pemrosesan, tetapi karena ini auto-commit, tampak agak kontradiktif bahwa yang boleh di-commit hanya item beberapa milidetik sebelum waktu auto-commit
polllalu memproses pesan secara durableTitik yang membingungkan adalah pemeriksaan auto-commit tidak terjadi secara asinkron setelah timeout, melainkan pada saat pemanggilan
pollberikutnyaKarena itu, penulisan seharusnya hanya bisa terjatuh jika sebelum memanggil
polllagi Anda tidak memproses pesan secara durable dan hanya menyimpannya, misalnya memakai pemrosesan asinkron, delay, atau queueIni adalah perilaku yang terdokumentasi pada library klien Java (https://kafka.apache.org/32/javadoc/org/apache/kafka/clients...), terlepas dari apakah implementasi saat ini benar-benar seperti itu
Protokol Kafka berada di antara level tinggi dan level rendah, sehingga keduanya tidak dijalankan dengan terlalu baik
Auto-commit adalah fitur level tinggi yang membantu mempermudah pembuatan aplikasi sederhana, tetapi jika tidak digunakan sesuai asumsi, tentu bisa gagal
Menurut saya, pengguna akhir sekarang sebaiknya tidak langsung memakai klien Kafka, melainkan implementasi level tinggi yang menangani detail-detail ini dengan benar. Untuk kebutuhan data, misalnya mesin pemrosesan stream; untuk kebutuhan aplikasi, misalnya mesin eksekusi berkelanjutan
Melihat halaman produk (https://buf.build/product/bufstream), saya penasaran bagaimana penjelasan “hanya berjalan di dalam AWS atau GCP VPC dan tidak menghubungi pihak luar” bisa selaras dengan penagihan berbasis penggunaan sebesar “$0.002 per GiB sebelum kompresi”
Rasanya mereka tidak mungkin menjalankan seluruh bisnis dengan sistem kejujuran begitu saja
Tentu ada risiko penyalahgunaan, tetapi itu bisa menjadi kompromi yang layak untuk menarik pelanggan tertentu
Jika source code-nya tidak dipublikasikan, saya sama sekali tidak akan percaya klaim “tidak menghubungi pihak luar”
Pernyataan “protokol transaksi Kafka pada dasarnya rusak dan harus direvisi” terdengar menyakitkan
Tetapi seperti biasa, investigasi dan tulisannya sangat bagus
Saya penasaran apakah Kyle pernah meninjau NATS JetStream. Ingin tahu pendapatnya
Beberapa orang menyarankan bahwa ini akan... bagaimana ya... menarik :-)
Saya tidak bisa menemukan proyek GitHub bufstream, jadi penasaran ada di mana
Hanya saja anehnya, ini juga tidak memiliki lisensi
Setelah membaca posting blog dan dokumen terkait, tampaknya “exactly-once delivery” Kafka didefinisikan sebagai properti dari pekerjaan read-process-write di mana worker membaca dari topik 1 lalu menulis ke topik 2, dan kedua topik itu berada dalam sistem Kafka logis yang sama
Jika benar, rasanya lebih tepat menyebut ini transaksi
Hanya saja ada dua cara memandang “exactly-once”
Salah satunya berarti efeknya tidak boleh terduplikasi atau hilang, seperti transaksi database
Yang lain lebih merupakan properti graf aliran data atas relasi pesan lintas topik-partisi, sedikit lebih dekat dengan konsistensi dalam ACID
Sama seperti sistem transaksi serializable dapat menjamin konsistensi tingkat domain tertentu, transaksi bisa digunakan untuk mencapai properti aliran data tersebut
Misalnya, serializability menjamin bahwa invariant yang terjaga saat melihat tiap transaksi secara terpisah juga tetap terjaga dalam riwayat eksekusi bersamaan
Bisa dibilang Kafka berusaha mencapai “exactly-once semantics” dengan cara seperti itu
Jangan sampai tertukar dengan https://www.warpstream.com/
Errata: “Transactions may observe none, part, or all” sepertinya seharusnya “Consumers may observe none, part, or all”
Semantik konsumen di luar transaksi lebih kabur
Semua pembacaan dalam workload ini terjadi dalam konteks transaksi, dan melewati jalur commit offset transaksional
Saya penasaran software ini dipakai untuk apa. Instrumentasi? Black box?
Tentu saja itu air mata kebahagiaan. Karena mendapat perhatian Jepsen sendiri sudah merupakan sebuah pencapaian