Maklumat dan konteks
Soalan kejuruteraan data dan pemprosesan strim ini disasarkan kepada jurutera data, jurutera platform penstriman dan peranan infrastruktur data bahagian belakang. Pengguna mungkin ranap atau mengalami imbangan semula (rebalance), dan pengeluar mungkin mencuba semula selepas kehilangan respons; hasil pemprosesan dihantar kembali ke Kafka terlebih dahulu. Jawapan mesti membezakan pemprosesan atomik Kafka daripada kesan sampingan luaran yang sewenang-wenangnya, bukannya menganggap exactly-once sebagai jaminan magis untuk keseluruhan sistem.
Perkara yang dinilai oleh penemu duga
- Sama ada anda boleh membahagikan exactly-once kepada keterlihatan output dan pemajuan atomik bagi ofset input.
- Sama ada anda boleh menggunakan pengeluar transaksi,
transactional.id,read_committeddan komit ofset manual dengan betul. - Sama ada anda boleh menerangkan pengguguran (abort), mula semula, pemagaran (fencing) dan pemulihan imbangan semula.
- Sama ada anda menyedari bahawa pangkalan data, enjin carian atau perkhidmatan HTTP memerlukan transaksinya sendiri, kunci keidempotenan atau penyesuaian (reconciliation).
Soalan penjelasan sebelum menjawab
Mula-mula, sahkan sama ada output kekal dalam Kafka, sama ada Kafka Streams digunakan, sama ada kumpulan pengguna membuat pelancaran berperingkat (rolling update) semasa penggunaan, dan sama ada sistem luaran mesti melakukan komit dalam transaksi yang sama seperti Kafka. Kemudian bezakan "satu output yang kelihatan bagi setiap input" daripada "satu kesan sampingan luaran"; perkara kedua memerlukan kerjasama daripada destinasi. Tanya tentang kependaman, saiz kelompok transaksi, tetingkap percubaan semula dan kelengahan (lag) yang boleh diterima.
Rangka kerja jawapan 30 saat
Saya akan meletakkan rekod input, output yang diubah dan ofset pengguna dalam satu transaksi Kafka. Pengeluar mendayakan transactional.id; pengguna menyahdayakan auto-komit dan menggunakan read_committed. Komit yang berjaya menjadikan ketiga-tiganya kelihatan bersama, manakala pengguguran menyembunyikan output dan meninggalkan ofset sebelum transaksi bermula. ID transaksi yang stabil dan unik akan memagar (fence) tikaan lama selepas mula semula. Ini merangkumi bacaan dan penulisan Kafka sahaja; pangkalan data luaran memerlukan hasil dan ofset dalam satu transaksi storan, atau keidempotenan, outbox dan penyesuaian.
Penyelesaian langkah demi langkah
1. Nyatakan sempadan jaminan
Reka bentuk Kafka menerangkan exactly-once topik-ke-topik sebagai pengemaskinian rekod output dan kedudukan pengguna secara atomik dalam satu transaksi. Ini tidak bermakna fungsi itu berjalan sekali sahaja atau permintaan HTTP sewenang-wenangnya tiba sekali sahaja. Nyatakan sempadan ini sebelum membincangkan konfigurasi.
2. Tulis output dan ofset secara atomik dengan pengeluar transaksi
Nyahdayakan auto-komit, proses satu kelompok, hantar output dan serahkan ofset kelompok tersebut sebagai sebahagian daripada transaksi. Pseudokod teras ialah:
producer.initTransactions();
while (running) {
ConsumerRecords<String, Order> records = consumer.poll(timeout);
producer.beginTransaction();
try {
for (ConsumerRecord<String, Order> r : records) {
producer.send(new ProducerRecord<>("orders-enriched", r.key(), enrich(r.value())));
}
producer.sendOffsetsToTransaction(offsets(records), groupMetadata);
producer.commitTransaction();
} catch (AbortableException e) {
producer.abortTransaction();
consumer.seekToCommitted();
}
}Melakukan komit output dan ofset secara bersama menghalang duplikasi yang kelihatan akibat kegagalan output-sebelum-ofset dan menghalang kehilangan akibat kegagalan ofset-sebelum-output. Kendalikan kelas pengecualian klien mengikut versi yang digunakan; mencuba semula setiap pengecualian secara membabi buta adalah tidak selamat.
3. Pastikan pengguna hanya melihat transaksi yang telah dikomit
Tetapkan isolation.level=read_committed dan kekalkan enable.auto.commit=false. read_uncommitted mendedahkan rekod daripada transaksi yang digugurkan, membolehkan pengguna menyebarkan output yang sepatutnya diundur balik (rolled back). read_committed menggunakan penanda transaksi untuk menyelaraskan keterlihatan dengan hasil komit.
4. Kendalikan mula semula, pemagaran dan imbangan semula
Berikan setiap tikaan pengguna yang aktif transactional.id yang unik dan stabil di seluruh kluster. Apabila tikaan baharu mendaftar dengan ID yang sama, Kafka menggugurkan transaksi aktif tikaan lama dan memagarnya (fences it), menghalang kedua-duanya daripada melakukan komit. Selepas pengguguran, aplikasi mesti mencipta semula atau mengundur balik kedudukan pengguna secara eksplisit dan memproses semula kelompok tersebut; ia tidak boleh diteruskan daripada kursor tempatan yang telah maju di dalam transaksi yang digugurkan. Penetapan sekatan (partition) memastikan seorang ahli kumpulan memiliki satu sekatan pada satu masa.
5. Terangkan sebab sistem luaran mesti bekerjasama
Jika hasil dihantar ke PostgreSQL, transaksi Kafka tidak merangkumi komit pangkalan data secara automatik. Ranap sistem selepas pangkalan data ditulis menyebabkan ofset Kafka tidak dikomit, jadi percubaan semula memerlukan kunci unik atau kemas kini versi bersyarat; komit ofset terlebih dahulu boleh menyebabkan penulisan hilang. Reka bentuk yang lebih kukuh meletakkan hasil dan ofset dalam satu transaksi pangkalan data, atau menulis rekod outbox/penyambung yang boleh dimainkan semula dan membiarkan destinasi melakukan penyahduplikasian dan penyesuaian. Tanpa kerjasama destinasi, janjikan pelaksanaan at-least-once dengan duplikasi yang boleh dikesan, bukan exactly-once hujung-ke-hujung.
6. Sahkan dengan suntikan kerosakan, bukan syot kilat konfigurasi
Lakukan simulasi ranap sebelum dan selepas sendOffsetsToTransaction, hilangkan respons komit, cetuskan imbangan semula, pagar tikaan lama, tamatkan tempoh transaksi dan jadikan sinki tidak tersedia. Gunakan output dengan read_committed dan semak kiraan output yang kelihatan bagi setiap ID input, ofset akhir, kelengahan selepas mula semula dan rekod transaksi yang digugurkan. Untuk sinki luaran, uji konflik kunci unik, main semula dan penyesuaian secara berasingan. Pantau kadar pengguguran, kelengahan pengguna, kependaman pemprosesan dan kiraan pemagaran.
Contoh jawapan berkualiti tinggi
Mula-mula, saya akan mengehadkan skop tuntutan tersebut. Jika kedua-dua input dan output berada dalam Kafka, saya akan menggunakan Kafka Streams atau gelung guna-ubah-hasil (consume-transform-produce) transaksi yang setaraf. Pengguna menyahdayakan auto-komit; pengeluar mempunyai transactional.id yang stabil dan unik; dan output bagi setiap kelompok berserta ofset dikomit dengan sendOffsetsToTransaction. Pengguna hiliran menggunakan read_committed, jadi output yang digugurkan tidak akan kelihatan. Semasa mula semula, ID transaksi yang sama akan memagar tikaan lama, dan pengguguran memerlukan pengunduran ke ofset komit yang terakhir. Jika destinasinya ialah PostgreSQL, saya tidak akan meluaskan jaminan Kafka menjadi tuntutan hujung-ke-hujung: saya akan meletakkan hasil dan ofset dalam satu transaksi pangkalan data, atau menggunakan outbox, kunci keidempotenan dan penyesuaian. Saya kemudiannya akan menyuntik ranap sistem, kehilangan respons, imbangan semula dan pemagaran serta mengesahkan output yang kelihatan, ofset dan keadaan destinasi bagi setiap ID input.
Kesilapan biasa
- Memanggil pengeluar idempoten sebagai exactly-once → ia hanya menghalang entri log pendua pada percubaan semula pengeluar dan tidak menjadikan output serta ofset bersifat atomik → tambah transaksi dan komit ofset.
- Hanya menetapkan
read_committed→ ia mengubah keterlihatan tetapi tidak melakukan komit transaksi atau memulihkan ranap sistem → laksanakan kitaran hayat transaksi penuh dengan auto-komit dimatikan. - Berkongsi satu
transactional.idmerentasi tikaan aktif → tikaan akan memagar antara satu sama lain dan menjejaskan daya pemprosesan → tetapkan ID unik yang stabil bagi setiap tikaan aktif. - Meneruskan daripada kedudukan tempatan selepas pengguguran → kedudukan tersebut mungkin berada di dalam kelompok yang belum dikomit → muat semula ofset yang dikomit atau lakukan carian (seek) secara eksplisit.
- Mendakwa Kafka menjadikan kesan sampingan pangkalan data berlaku sekali sahaja → kedua-dua sistem tidak mempunyai komit atomik automatik → gunakan transaksi destinasi, penulisan idempoten, outbox atau penyesuaian.
Soalan susulan dan jawapan
Adakah exactly-once bermaksud fungsi perniagaan dijalankan sekali sahaja?
Tidak. Fungsi ini boleh dijalankan dan kemudian dijalankan semula selepas kegagalan komit; jaminannya ialah output Kafka yang kelihatan dan komit ofset adalah selaras. Pastikan fungsi bebas daripada kesan sampingan luaran jika boleh, atau jadikan kesan tersebut boleh diulang dan disesuaikan.
Mengapakah read_committed masih boleh menambah kependaman?
Pengguna mesti melangkau rekod yang digugurkan dan menunggu transaksi selesai. Transaksi yang terbuka atau berdurasi panjang meningkatkan kependaman keterlihatan dan kelengahan. Hadkan saiz kelompok dan had masa tamat (timeout), pantau tempoh transaksi dan pertimbangkan kosnya berbanding pemprosesan at-least-once untuk laluan kependaman rendah.
Apakah yang berlaku apabila imbangan semula berlaku di pertengahan transaksi?
Pemilik baharu menyambung semula daripada ofset terakhir yang dikomit. Transaksi semasa harus digugurkan; tikaan lama dilepaskan atau dipagar; dan tikaan baharu memproses semula kelompok yang belum dikomit. Jangan komit ofset melainkan ahli semasa masih memiliki sekatan tersebut.
Bagaimanakah anda akan menghantar output Kafka ke PostgreSQL?
Tulis hasil perniagaan dan kedudukan pengguna dalam satu transaksi PostgreSQL, atau tulis baris outbox dengan ID peristiwa unik dan biarkan geganti yang boleh dipercayai menghantarnya. Jika kedua-duanya tidak boleh berkongsi transaksi, pilih penghantaran at-least-once, kekangan unik, kemas kini versi bersyarat dan penyesuaian; jangan promosikannya sebagai exactly-once.
Metrik manakah yang membuktikan bahawa reka bentuk itu berfungsi?
Kira output unik read_committed yang kelihatan bagi setiap ID peristiwa input, kemudian hubung kaitkan ofset yang dikomit, kiraan pengguguran dan pemagaran, kelengahan pengguna, volum main semula selepas mula semula dan konflik kunci pendua pada sinki. Simpan sampel sebelum dan selepas suntikan kerosakan; syot kilat konfigurasi tidak dapat membuktikan keterpeliharaan fungsi tersebut.