7.6

View in English

7.6 Data real-time dan streaming

Tinjauan dan motivasi

Sebagian besar yang Anda ketahui tentang pipeline data mengasumsikan data diam. Anda mengumpulkan catatan sehari, menjalankan pekerjaan semalaman, dan membaca hasilnya di pagi hari. Data real-time dan streaming membalik asumsi itu. Alih-alih memproses tumpukan data yang sudah selesai, Anda memproses aliran peristiwa tak berujung saat tiba, dan menghasilkan jawaban secara terus-menerus. Inilah beda antara pemrosesan batch, yang beroperasi pada dataset terbatas dan lengkap, dan pemrosesan stream, yang beroperasi pada aliran tak terbatas dan tak pernah selesai.

Bagi tim besar, streaming muncul begitu latensi mulai penting bagi bisnis. Keputusan penipuan yang tiba sejam terlambat tak berharga. Sinyal personalisasi yang mendarat besok tidak mempersonalisasi apa pun. Dasbor operasional yang tertinggal dari kenyataan sepanjang satu shift menyesatkan orang yang mengawasinya. Bab 7.2 (rekayasa data) berargumen bahwa Anda harus memilih batch sebagai bawaan dan meraih streaming hanya di mana latensi benar-benar terbayar, dan bab ini membawa Anda sisanya: kapan real-time layak biayanya, dan bagaimana membangunnya tanpa membakar anggaran operasi Anda. Streaming berdekatan dengan pola olah pesan berbasis peristiwa di bab 3.12 (arsitektur berbasis peristiwa dan pesan), pilihan penyimpanan di bab 3.4 (arsitektur data dan penyimpanan), dan praktik telemetri di bab 9.2 (observabilitas dan telemetri).

Pengaturan enterprise dan pemerintah menaikkan taruhan. Sebuah bank menilai transaksi untuk penipuan dalam waktu yang dibutuhkan pembaca kartu untuk berkedip. Sebuah dinas transit melacak kendaraan dan memprediksi kedatangan untuk jutaan penumpang. Sebuah lembaga tunjangan mengawasi anomali dalam klaim sambil menjaga catatan keputusan yang dapat diaudit. Dalam semua ini, nilai datang dari bertindak atas data selagi masih segar, dan risiko datang dari bertindak atas data yang salah, tidak lengkap, atau mustahil direkonstruksi kemudian. Bab ini beropini tentang keduanya.

Prinsip utama

  • Raih streaming hanya ketika latensi punya nilai bisnis yang jelas; batch lebih murah dan sederhana.
  • Bedakan data terbatas (finit) dari data tak terbatas (tak berujung), dan rancang sesuai.
  • Perlakukan waktu peristiwa, bukan waktu kedatangan, sebagai sumber kebenaran, dan rencanakan untuk data terlambat dan tak berurutan.
  • Window dan watermark adalah cara mendapatkan jawaban finit dari stream tak berhingga.
  • Pilih hasil effectively-once lewat sink idempoten daripada janji exactly-once yang rapuh.
  • Pemrosesan berstatus butuh checkpoint agar dapat pulih tanpa kehilangan atau menghitung ganda.
  • Rancang untuk backpressure dan pemrosesan ulang sejak hari pertama, bukan renungan belakangan.
  • Jaga logika streaming dapat diamati dan diaudit; stream senyap lebih buruk daripada batch yang gagal.

Rekomendasi

Benarkan real-time sebelum membangunnya

Keputusan streaming terpenting adalah apakah akan streaming sama sekali. Real-time kira-kira menggandakan kompleksitas dan biaya operasional Anda, karena Anda menukar pekerjaan yang berjalan lalu berhenti dengan sistem yang harus sehat setiap detik. Sebelum berkomitmen, namai keputusan yang dimungkinkan data segar dan biaya keputusan itu tiba terlambat. Penilaian penipuan, peringatan operasional, dan personalisasi langsung biasanya melewati batas. Dasbor yang dilihat manusia dua kali sehari nyaris tidak pernah, seberapa pun memuaskannya “real-time” terdengar dalam rapat perencanaan. Tuliskan persyaratan latensi sebagai angka, dalam detik atau menit, dan periksa terhadap kenyataan. Banyak yang disebut orang real-time terlayani baik oleh micro-batch yang berjalan setiap beberapa menit dengan sebagian kecil biayanya.

Rancang berdasarkan waktu peristiwa, bukan waktu pemrosesan

Gagasan tersulit dalam streaming adalah bahwa peristiwa terjadi pada satu saat dan diproses pada saat lain. Waktu peristiwa adalah kapan hal itu benar-benar terjadi, misalnya ketika penumpang menempelkan kartu. Waktu pemrosesan adalah kapan sistem Anda sempat menanganinya. Keduanya terus menyimpang: ponsel kehilangan sinyal di terowongan dan mengunggah tiga menit tap sekaligus, gangguan jaringan mengacak ulang pesan, partisi tertinggal. Jika Anda menghitung berdasarkan waktu pemrosesan, angka Anda bergoyang mengikuti infrastruktur alih-alih mencerminkan dunia. Masalah terlambat dan tak berurutan ini inti disiplin, dan terhubung langsung dengan pemodelan peristiwa dalam arsitektur berbasis peristiwa. Cap setiap peristiwa dengan waktu peristiwanya di sumber, bawa cap waktu itu melalui seluruh pipeline, dan hitung hasil Anda terhadapnya.

Gunakan window dan watermark untuk mendapatkan jawaban finit

Stream tak terbatas tak pernah berakhir, jadi “hitung peristiwa” tak punya jawaban sampai Anda membatasinya. Window melakukan pembatasan itu. Window tumbling memotong waktu menjadi keranjang tetap tak tumpang tindih, misalnya setiap menit. Window sliding tumpang tindih, sehingga window lima menit yang maju setiap menit memberi Anda angka bergerak yang mulus. Window session mengelompokkan ledakan aktivitas yang dipisahkan jeda tak aktif, yang cocok untuk sesi pengguna. Setelah punya window, Anda perlu memutuskan kapan window selesai, karena data terlambat mungkin masih tiba. Watermark adalah estimasi sistem bahwa ia mungkin telah melihat semua peristiwa hingga waktu peristiwa tertentu. Ketika watermark melewati akhir window, Anda memancarkan hasilnya. Setel berapa lama Anda menunggu: tahan window lebih lama dan Anda menoleransi lebih banyak keterlambatan dengan biaya latensi dan memori, tutup lebih cepat dan Anda berisiko menjatuhkan yang tertinggal. Putuskan secara eksplisit apa yang terjadi pada data yang tiba setelah window tertutup, entah dibuang, dicatat, atau dipancarkan sebagai koreksi.

Jadikan sink idempoten dan pilih effectively-once

Jaminan pengiriman terdengar sederhana dan tidak. Pengiriman at-least-once berarti setiap peristiwa diproses, tetapi sebagian mungkin diproses lebih dari sekali setelah percobaan ulang, sehingga hitungan dapat menggelembung. Exactly-once terdengar ideal tetapi mahal dan, diambil harfiah lintas sistem eksternal sembarang, sering mustahil. Target praktisnya adalah effectively-once: hasil yang teramati seolah setiap peristiwa diproses sekali, meski mesin di bawahnya mencoba ulang. Anda sampai di sana dengan membuat sink idempoten aman untuk ditulisi berulang kali, memakai kunci deterministik dan upsert agar peristiwa yang diputar ulang menimpa alih-alih menduplikasi. Padukan pengiriman at-least-once dengan penulisan idempoten dan Anda mendapat hasil benar tanpa membayar koordinasi transaksional berat di mana-mana. Sisihkan mesin exactly-once sejati untuk tempat sempit yang benar-benar membutuhkannya.

Checkpoint pemrosesan berstatus agar dapat pulih

Banyak komputasi streaming berguna bersifat berstatus: hitungan berjalan, join lintas stream, deduplikasi, model penipuan yang mengingat perilaku terbaru. Status itu hidup di memori dan akan lenyap ketika proses dimulai ulang. Checkpointing secara berkala memotret status dan posisi stream bersama-sama, sehingga setelah crash sistem melanjutkan dari titik konsisten alih-alih memutar ulang segalanya atau kehilangan ingatannya. Tentukan ukuran status dengan sengaja, karena status tak terbatas adalah cara umum menghabiskan memori pekerjaan streaming di produksi. Pakai kedaluwarsa dan time-to-live pada status yang tak lagi Anda butuhkan, dan pantau ukuran status sebagai metrik kelas satu. Waktu pemulihan setelah kegagalan adalah perhatian service-level nyata, jadi uji sebelum pengguna Anda melakukannya.

Streaming dari basis data operasional dengan change data capture

Anda sering ingin bereaksi terhadap perubahan di basis data yang tak pernah dirancang memancarkan peristiwa. Change data capture (CDC) menyelesaikan ini dengan membaca log transaksi basis data dan mengubah setiap insert, update, dan delete menjadi stream peristiwa perubahan. Ini jauh lebih baik daripada melakukan polling tabel pada timer, yang lambat, melewatkan keadaan antara, dan menghantam sumber. CDC memungkinkan Anda menjaga indeks pencarian, cache, penyimpanan analitik, atau layanan hilir tetap tersinkron terus-menerus dengan sistem pencatat, dan melakukannya tanpa perubahan invasif pada aplikasi. Perlakukan stream perubahan sebagai produk data kelas satu: versikan skemanya, dokumentasikan maknanya, dan awasi lag-nya, karena segala yang di hilir mewarisi lag itu.

Pilih arsitektur streaming-first daripada memelihara dua basis kode

Arsitektur Lambda klasik menjalankan lapisan batch untuk riwayat akurat dan lengkap di samping lapisan speed untuk hasil segar dan perkiraan, lalu menggabungkannya. Itu berfungsi, tetapi memaksa Anda menulis dan memelihara logika bisnis yang sama dua kali, di dua sistem, dan merekonsiliasi perbedaannya selamanya. Arsitektur Kappa meruntuhkan ini: simpan log peristiwa yang tahan lama dan dapat diputar ulang dan jalankan semua pemrosesan sebagai pemrosesan stream, memproses ulang riwayat dengan memutar ulang log ketika logika berubah. Industri hanyut ke bentuk streaming-first ini karena satu basis kode jauh lebih murah dipelihara dan dinalar. Jika Anda dapat mengekspresikan kebutuhan batch sebagai pemutaran ulang atas log peristiwa yang disimpan, Anda menghindari pajak dua-basis-kode sepenuhnya. Pakai broker berbasis log yang menyimpan riwayat agar pemrosesan ulang berupa memutar mundur, bukan membangun ulang.

Ekspos stream sebagai SQL, materialized view, dan OLAP real-time

Tidak semua orang yang butuh streaming harus menulis kode pemrosesan stream tingkat rendah. Streaming SQL memungkinkan analis dan insinyur mengekspresikan window, join, dan agregasi dalam bahasa yang sudah mereka kenal, dan menjaga hasil tetap mutakhir terus-menerus sebagai materialized view. Untuk kueri analitis berlatensi rendah atas data segar, penyimpanan online analytical processing (OLAP) real-time mengingesti stream dan menjawab kueri iris-dan-potong dalam milidetik, yang menggerakkan dasbor operasional yang benar-benar langsung. Padukan ini dengan praktik analitik produk di bab 7.4 (analitik produk dan eksperimen) ketika tujuannya umpan balik cepat atas fitur dan eksperimen. Pilih perkakas tingkat lebih tinggi ini di mana cocok, dan sisihkan pemroses stream tulisan tangan untuk logika yang tak dapat mereka ekspresikan.

Rencanakan backpressure dan pemrosesan ulang sejak awal

Stream dapat tiba lebih cepat daripada yang dapat Anda proses. Backpressure adalah mekanisme yang memungkinkan konsumen lambat memberi sinyal ke hulu untuk melambat alih-alih tumbang atau menjatuhkan data secara senyap. Pastikan setiap tahap dalam pipeline Anda menghormatinya, dan pantau lag konsumen sebagai metrik utama, karena lag yang tumbuh adalah peringatan paling awal bahwa Anda kalah dalam perlombaan. Pemrosesan ulang adalah kemampuan lain yang diharapkan orang telah mereka bangun. Ketika Anda menemukan bug atau mengubah aturan, Anda ingin memutar ulang riwayat melalui logika yang dikoreksi. Itu hanya mungkin jika log peristiwa Anda menyimpan cukup riwayat dan sink Anda cukup idempoten untuk menyerap pemutaran ulang. Rancang keduanya sejak hari pertama; memasangnya belakangan di bawah tekanan insiden menyengsarakan.

Trade-off: kelebihan dan kekurangan

PilihanKelebihanKekuranganPaling cocok
BatchSederhana, murah, mudah diuji dan di-backfillLatensi tinggi, basi di antara prosesPelaporan, sebagian besar analitik
Micro-batch (menit)Hampir real-time, jauh lebih sederhana daripada streamingTidak benar-benar instanDasbor “real-time”
Streaming sejati (sub-detik)Reaksi instan, hasil terus-menerusKompleks, mahal, sulit diujiPenipuan, peringatan, personalisasi langsung
At-least-once + sink idempotenHasil benar, terjangkau, tangguhMembutuhkan desain kunci yang disiplinSebagian besar pipeline streaming
Mesin exactly-onceJaminan kuat ujung ke ujungMahal, terbatas lintas sistemJalur sempit berisiko tinggi
Lambda (batch + speed)Riwayat akurat plus pandangan segarDua basis kode untuk dipeliharaMigrasi warisan
Kappa (streaming-first)Satu basis kode, dapat diputar ulangMembutuhkan log tahan lama yang disimpanPlatform streaming baru

Ketegangan sentralnya adalah latensi melawan kompleksitas. Setiap langkah menuju real-time berbiaya beban operasional, kesulitan pengujian, dan uang, dan imbalannya tidak linear: berpindah dari harian ke setiap beberapa menit murah dan sering cukup, sementara dari menit ke sub-detik adalah tempat biaya terkonsentrasi. Selesaikan ketegangan dengan menghargai keputusan, bukan teknologi. Tanyakan tindakan apa yang dimungkinkan kesegaran dan berapa biaya keterlambatan, lalu beli hanya sebanyak pengurangan latensi yang dibenarkan tindakan itu. Ketika Anda memang butuh streaming, bersandarlah pada pengiriman at-least-once dengan sink idempoten dan log streaming-first, karena kombinasi itu memberi kebenaran dan kemampuan diputar ulang tanpa jaminan terberat.

Pertanyaan untuk didiskusikan dengan tim Anda

  1. Keputusan apa yang sebenarnya dimungkinkan data real-time bagi kita, dan berapa biayanya ketika data itu tiba semenit terlambat alih-alih seketika? Ini pertanyaan yang harus menggerbangi setiap proyek streaming, karena streaming kira-kira menggandakan biaya dan kompleksitas operasional Anda dibanding batch. Tim besar dapat membakar kuartal-kuartal membangun platform real-time yang melayani dasbor yang dicek manusia dua kali sehari, yaitu uang yang dibakar. Bawa tindakan konkret yang digerakkan data, entah memblokir transaksi curang, memanggil operator, atau mengubah apa yang dilihat pengguna, dan beri angka pada biaya latensi untuk masing-masing. Jika jawaban jujurnya bahwa micro-batch lima menit akan memenuhi kebutuhan, itu temuan yang layak dirayakan, bukan disembunyikan. Jawabannya harus langsung mengubah apakah Anda membangun streaming sejati, puas dengan micro-batch, atau tetap di batch.

  2. Bagaimana kita menangani peristiwa terlambat dan tak berurutan, dan apa yang terjadi pada data yang tiba setelah window tertutup? Data terlambat dan tak berurutan adalah bagian sulit streaming, dan tim yang melewatkan pertanyaan ini menemukannya di produksi ketika angka mereka menolak berekonsiliasi. Tekanan yang bersaing adalah latensi dan kebenaran: tahan window lebih lama untuk menangkap yang tertinggal dan Anda menunda setiap hasil dan mengonsumsi lebih banyak memori, tutup lebih cepat dan Anda diam-diam menjatuhkan data nyata. Bawa bukti tentang seberapa terlambat data Anda sebenarnya tiba, diukur sebagai jarak antara waktu peristiwa dan waktu pemrosesan di seluruh sumber Anda, karena sumber seluler di terowongan berperilaku sangat berbeda dari peristiwa sisi server. Putuskan secara eksplisit apakah data terlambat dibuang, dicatat, atau memicu koreksi, dan pastikan semua orang di hilir tahu yang mana. Dalam konteks pemerintah di mana angka harus dapat dipertahankan, diam-diam menjatuhkan peristiwa terlambat dapat menjadi masalah kepatuhan, jadi kebijakannya harus disengaja dan terdokumentasi.

  3. Apakah sink kita cukup idempoten sehingga kita dapat dengan aman memutar ulang riwayat, dan apakah log peristiwa kita menyimpan cukup untuk memungkinkan pemutaran ulang? Pemrosesan ulang adalah kemampuan yang paling sering diharapkan tim telah dibangun dan paling sering tidak, dan ia bergantung pada dua hal yang bekerja bersama: sink idempoten yang menyerap peristiwa terputar ulang tanpa menduplikasi, dan log tahan lama yang menyimpan cukup riwayat untuk diputar ulang. Tanpa keduanya, memperbaiki bug logika berarti Anda tak dapat menghitung ulang periode terdampak dengan bersih, dan Anda terjebak menambal angka dengan tangan di bawah tekanan. Bawa jendela retensi Anda saat ini dan uji konkret: pilih bug nyata dari kuartal lalu dan tanyakan apakah Anda dapat memutar ulang logika yang dikoreksi atas data terdampak. Tarikan melawannya adalah biaya, karena menyimpan riwayat dan merancang penulisan idempoten memakan penyimpanan dan disiplin di muka. Tetapi alternatifnya muncul pada saat terburuk, selama insiden, jadi jawabannya membentuk seberapa banyak Anda berinvestasi pada kemampuan diputar ulang sebelum Anda membutuhkannya.

  4. Ketika pekerjaan streaming crash, secepat apa ia harus pulih, berapa banyak status yang boleh dipegangnya, dan sudahkah kita benar-benar mengukur waktu pemulihan di bawah beban produksi? Pekerjaan batch yang mati dapat dijalankan ulang besok, tetapi stream yang selalu hidup dan mati adalah pemadaman yang sedang berlangsung, dan pekerjaan berstatus yang memegang hitungan berjalan, join, atau model penipuan dapat kehilangan menit-menit memori atau butuh waktu lama memuat ulang status setelah restart. Bagi tim besar, di sinilah detail tak glamor diam-diam menetapkan ketersediaan nyata Anda: status tak terbatas tumbuh sampai pekerjaan kehabisan memori, dan pemulihan checkpoint yang lambat mengubah gangguan sepuluh detik menjadi sepuluh menit. Tekanan yang bersaing adalah kesegaran melawan keamanan, karena checkpoint lebih sering memperpendek pemulihan tetapi menambah overhead, dan retensi status yang longgar memperbaiki akurasi tetapi berisiko kehabisan memori. Bawa recovery time objective konkret, ukuran status Anda saat ini dan kurva pertumbuhannya, interval checkpoint Anda, dan hasil latihan failover nyata alih-alih perkiraan penuh harapan. Dalam pengaturan enterprise dan pemerintah di mana stream mendukung penilaian penipuan atau umpan keselamatan publik, jalur pemulihan yang belum diuji adalah risiko operasional yang Anda terima tanpa mengukur, jadi perlakukan latihan sebagai persyaratan, bukan bagus-jika-ada.

  5. Apakah kita menjalankan satu basis kode streaming-first atau lapisan batch dan lapisan speed terpisah, dan berapa sebenarnya biaya menjaga keduanya terekonsiliasi? Pola Lambda berupa lapisan batch untuk riwayat akurat plus lapisan speed untuk hasil segar memaksa Anda menulis logika bisnis yang sama dua kali, di dua sistem, dan merekonsiliasi jawabannya selamanya, sedangkan bentuk streaming-first (Kappa) menjaga log tahan lama yang dapat diputar ulang dan menjalankan semua pemrosesan sebagai pemrosesan stream. Bagi organisasi besar, logika terduplikasi adalah tempat penyimpangan dan angka yang disengketakan berkembang biak, karena aturan berubah di satu lapisan dan tidak di yang lain, dan insinyur menghabiskan waktu nyata menjelaskan mengapa keduanya tidak sepakat. Tarikan untuk mempertahankan keduanya adalah inersia dan kenyamanan lapisan batch yang teruji, jadi timbang itu terhadap pajak pemeliharaan dengan jujur. Bawa daftar komputasi yang saat ini Anda jalankan di kedua tempat, insiden yang disebabkan kedua lapisan tidak sepakat, dan penilaian apakah log peristiwa Anda menyimpan cukup riwayat untuk mengekspresikan kebutuhan batch sebagai pemutaran ulang. Dalam konteks pemerintah dan enterprise teraudit, dua lapisan yang dapat melaporkan angka berbeda untuk periode yang sama adalah liabilitas kepatuhan tersendiri, karena Anda harus dapat mengatakan angka mana yang otoritatif dan mengapa.

  6. Siapa yang mengoperasikan sistem selalu-hidup ini ketika rusak pukul tiga pagi, dan sudahkah kita menganggarkan beban on-call dan keterampilan spesialis yang dituntutnya, atau kita mengasumsikan staf berbentuk batch? Streaming menggeser biaya dari bangun ke jalankan: sistem harus sehat setiap detik, yang berarti cakupan on-call nyata, insinyur fasih dalam waktu peristiwa, watermark, status, dan semantik pengiriman, dan pengujian yang lebih sulit daripada untuk pekerjaan yang berjalan lalu berhenti. Tim rutin menyetujui platform streaming atas dasar kemampuannya dan tak pernah mendanai orang yang menjaganya tetap hidup, sehingga platform merosot dan kepercayaan terkikis. Trade-off-nya adalah cakupan melawan keberlanjutan: setiap pipeline real-time tambahan adalah satu hal lagi yang dapat memanggil seseorang, jadi pertanyaannya apakah latensi yang dibelinya membenarkan komitmen operasional permanen. Bawa inventaris jujur siapa yang memiliki setiap stream di produksi, rotasi on-call Anda saat ini dan ruang geraknya, dan di mana keahlian waktu-peristiwa sebenarnya berada, entah perekrutan, mitra, atau layanan terkelola. Untuk badan publik atau enterprise besar, tambahkan lead time pengadaan dan perekrutan dan opsi layanan terkelola mana pun, karena platform real-time yang bergantung pada talenta langka yang tak dapat Anda rekrut atau pertahankan adalah rencana menjalankan sistem rawan pemadaman dengan staf kurang.

Lensa sektor

Startup. Streaming jarang langkah pertama Anda, dan mendirikan platform berat dapat menenggelamkan tim kecil. Pilih satu sinyal yang menyentuh nilai inti Anda, taruh peristiwa di satu broker berbasis log yang menyimpan, dan jalankan pemroses ringan dengan sink idempoten berkunci agar percobaan ulang at-least-once tak pernah menghitung ganda. Simpan riwayat beberapa hari agar Anda dapat memutar ulang melalui logika yang diperbaiki, dan pilih layanan streaming terkelola daripada mengoperasikan klaster sendiri, karena sumber daya Anda yang paling langka adalah perhatian rekayasa.

Bisnis kecil. Anda mungkin tidak punya spesialis streaming dan tak berselera menjalankan infrastruktur selalu-hidup, jadi perlakukan real-time sebagai sesuatu yang Anda beli di dalam perkakas yang sudah Anda pakai, bukan sistem yang Anda isi stafnya. Bingkai kebutuhan sebagai pertanyaan latensi dengan angka terlampir, dan dalam banyak kasus micro-batch yang menyegarkan setiap beberapa menit akan memenuhinya dengan sebagian kecil biaya dan risiko. Pilih vendor yang fitur real-time-nya transparan tentang lag dan mudah dimundurkan, dan sisihkan streaming pesanan untuk kasus langka di mana data segar langsung menggerakkan pendapatan atau keselamatan.

Enterprise. Masalahnya konsistensi dan biaya lintas banyak tim: platform berbasis log bersama, kebijakan waktu-peristiwa dan data terlambat standar, dan sink idempoten agar kelompok berhenti menciptakan ulang pipeline rapuh. Anggarkan beban operasi selalu-hidup dan on-call secara eksplisit, bakukan log streaming-first agar menghindari basis kode batch terduplikasi, dan kelola stream sebagai produk data teratur dengan pemilik, versioning skema, dan lag terpantau alih-alih sebaran pekerjaan pesanan. Lacak latensi, waktu pemulihan, dan biaya per stream sebagai metrik portofolio.

Pemerintah. Kemampuan diaudit dan akuntabilitas publik membentuk setiap pilihan. Simpan setiap peristiwa yang diproses dalam log tahan lama agar angka yang dilaporkan kepada badan pengawas, jumlah penumpang, anomali tunjangan, keputusan penipuan, dapat direkonstruksi persis, dan jadikan kebijakan data terlambat eksplisit dan terdokumentasi alih-alih diam-diam menjatuhkan peristiwa. Pengadaan harus menuntut portabilitas data dan pengungkapan jaminan pengiriman dan retensi layanan terkelola, dan setiap restatement setelah perubahan aturan harus berupa pemutaran ulang yang dapat dipertahankan melalui logika yang dikoreksi, bukan tambalan manual yang tak dapat ditelusuri siapa pun.

Contoh

Startup. Aplikasi konsumen ingin menampilkan feed aktivitas langsung kepada pengguna dan menandai login mencurigakan saat terjadi. Tim menolak mendirikan platform streaming berat. Mereka menaruh peristiwa di satu broker berbasis log yang menyimpan, menjalankan pemroses stream ringan untuk logika risiko login, dan memberi makan penyimpanan OLAP real-time yang menggerakkan feed aktivitas. Setiap sink berkunci dan idempoten, sehingga percobaan ulang at-least-once tak pernah menghitung ganda. Ketika kemudian menemukan bug dalam aturan risiko, mereka cukup memutar ulang log melalui logika yang diperbaiki semalaman, karena mereka menyimpan riwayat seminggu dan tak pernah membutuhkan basis kode batch kedua.

Enterprise. Sebuah bank ritel menilai setiap transaksi kartu untuk penipuan dalam jendela otorisasi, menggabungkan stream transaksi langsung dengan model berstatus perilaku akun terbaru. Checkpointing memungkinkan layanan penilaian pulih dari kegagalan node dalam hitungan detik tanpa kehilangan ingatan beberapa menit terakhir. Secara terpisah, change data capture mengalirkan pembaruan dari basis data perbankan inti ke indeks pencarian dan layanan personalisasi, menjaga keduanya segar tanpa polling. Dasbor operasional membaca dari penyimpanan OLAP real-time agar tim risiko dan operasi menyaksikan bisnis bergerak, dan seluruh pipeline memancarkan telemetri lag dan throughput yang dijelaskan di bab 9.2.

Pemerintah. Sebuah otoritas transit metropolitan mengingesti posisi kendaraan dan tap tarif untuk memprediksi kedatangan dan memantau kepadatan secara real-time, memberi makan aplikasi publik dan pusat operasi. Karena penumpang di terowongan mengunggah tap dalam ledakan tertunda, tim menghitung jumlah penumpang berdasarkan waktu peristiwa dengan watermark yang disetel pada keterlambatan teramati, dan mencatat peristiwa yang tiba setelah window-nya tertutup alih-alih menjatuhkannya secara senyap. Setiap peristiwa yang diproses disimpan dalam log teraudit sehingga angka penumpang yang dilaporkan kepada badan pengawas dapat direkonstruksi persis. Ketika aturan tarif berubah, mereka memutar ulang periode terdampak melalui logika yang dikoreksi dan menghasilkan restatement yang dapat dipertahankan.

Kasus bisnis: motivasi, ROI, dan TCO

Imbal hasil data real-time datang dari bertindak selagi tindakan masih berarti. Penipuan yang tertangkap selama otorisasi mencegah kerugian yang hanya akan dilaporkan batch malam. Personalisasi yang merespons dalam sesi menaikkan konversi dengan cara yang tak dapat dicapai rekomendasi besok. Pemantauan operasional yang mencerminkan masa kini memungkinkan Anda mengintervensi sebelum masalah kecil menjadi pemadaman atau insiden publik. Dalam setiap kasus, nilainya adalah selisih antara bertindak sekarang dan bertindak nanti, dan selisih itulah yang harus Anda kuantifikasi ketika mengajukan kasus.

Total biaya kepemilikan lebih tinggi daripada batch, dan kejujuran tentang itu melindungi kredibilitas Anda. Anda membayar infrastruktur selalu-hidup, insinyur yang memahami waktu peristiwa, watermark, status, dan semantik pengiriman, dan beban pengujian serta on-call yang lebih sulit dari sistem yang harus sehat terus-menerus alih-alih berjalan lalu berhenti. Arsitektur streaming-first di atas log yang disimpan menurunkan biaya berjalan dengan menyelamatkan Anda dari basis kode batch duplikat, dan memilih at-least-once dengan sink idempoten menghindari biaya mesin exactly-once ujung ke ujung. Kesalahan termahal adalah membangun real-time di mana micro-batch atau batch sudah cukup, jadi argumen biaya terkuat sering keputusan untuk tidak streaming. Bingkai pitch kepada pimpinan di sekitar keputusan peka-latensi spesifik dan imbal hasil terukurnya, dan sama jelasnya tentang di mana tetap di batch menghemat uang tanpa kehilangan nilai.

Anti-pola dan jebakan

  • Membangun streaming demi gengsi ketika micro-batch setiap beberapa menit akan memenuhi kebutuhan.
  • Menghitung berdasarkan waktu pemrosesan, sehingga angka Anda bergoyang mengikuti infrastruktur alih-alih dunia.
  • Mengabaikan data terlambat dan tak berurutan sampai rekonsiliasi gagal di produksi.
  • Mengejar exactly-once harfiah di mana-mana alih-alih at-least-once dengan sink idempoten.
  • Status tak terbatas tanpa kedaluwarsa, diam-diam tumbuh sampai pekerjaan kehabisan memori.
  • Tanpa checkpointing, sehingga restart kehilangan status atau memaksa pemutaran ulang penuh.
  • Mem-polling basis data operasional pada timer alih-alih memakai change data capture.
  • Memelihara lapisan batch dan speed Lambda dengan logika terduplikasi yang menyimpang.
  • Jendela retensi terlalu pendek untuk memutar ulang riwayat ketika Anda menemukan bug.
  • Stream tanpa metrik lag, throughput, atau kesegaran, gagal secara senyap.

Model kematangan

  • Tingkat 1, Memulai: Segalanya batch, atau beberapa pekerjaan streaming buatan tangan berjalan reaktif tanpa pemantauan. Angka dihitung berdasarkan waktu pemrosesan, data terlambat diabaikan, dan restart kehilangan status. Tak ada yang dapat memutar ulang riwayat untuk memperbaiki bug, dan masalah ditemukan ketika angka hilir menolak berekonsiliasi.
  • Tingkat 2, Mengembangkan: Sebagian tim menjalankan pipeline streaming inti di broker berbasis log dengan checkpointing, dan mereka membedakan waktu peristiwa dari waktu pemrosesan dan memakai window dasar. Praktik tidak konsisten dari tim ke tim: pengiriman at-least-once tetapi tidak semua sink idempoten, penanganan data terlambat diimprovisasi, dan lag diawasi secara informal alih-alih diberi peringatan.
  • Tingkat 3, Membakukan: Waktu peristiwa, watermark, dan kebijakan data terlambat eksplisit didokumentasikan dan diterapkan di seluruh organisasi. Sink idempoten untuk hasil effectively-once, status punya kedaluwarsa, dan change data capture memberi makan sistem hilir menurut konvensi. Log yang disimpan mendukung pemutaran ulang, dan lag, throughput, serta kesegaran dipantau dengan peringatan sebagai standar seluruh organisasi alih-alih kebiasaan per tim.
  • Tingkat 4, Mengelola: Properti streaming diukur dan dikendalikan terhadap garis dasar. Setiap pipeline membawa service-level objective untuk latensi ujung ke ujung, lag konsumen, waktu pemulihan, kemiringan waktu-peristiwa, tingkat peristiwa terlambat, ukuran status, dan biaya per juta peristiwa, semuanya dilacak terhadap target yang disepakati dan memberi peringatan atas regresi. Pemulihan dilatih dan diukur waktunya alih-alih diasumsikan, ruang backpressure dan pertumbuhan status diawasi sebagai sinyal kapasitas, dan stream baru harus melewati metrik ini sebelum ke produksi.
  • Tingkat 5, Mengorkestrasi: Arsitektur streaming-first melayani kebutuhan segar dan historis dari satu log yang dapat diputar ulang, dan streaming SQL, materialized view, serta OLAP real-time membuat data segar dapat diakses luas. Pemrosesan ulang rutin dan teruji, platform menskalakan otomatis dan menyeimbangkan ulang terhadap beban dan biaya terukur, dan stream dipensiunkan, ditentukan ulang cakupannya, atau diganti atas bukti. Streaming terintegrasi dengan perencanaan bisnis dan risiko, dan setiap stream dapat diamati dan diaudit ujung ke ujung seiring gambaran beban dan biaya bergeser.

Gagasan untuk didiskusikan

  1. Di mana dalam tumpukan Anda “real-time” benar-benar layak biayanya, dan di mana ia keinginan tak diperiksa?
  2. Seberapa besar jarak antara waktu peristiwa dan waktu pemrosesan di seluruh sumber Anda, dan apakah Anda mengukurnya?
  3. Dapatkah Anda meruntuhkan pengaturan batch-dan-speed Lambda menjadi satu basis kode streaming-first, dan apa yang akan menghalangi?
  4. Sink Anda yang mana benar-benar idempoten, dan dapatkah Anda dengan aman memutar ulang data kuartal lalu melalui logika yang dikoreksi hari ini?
  5. Apa kebijakan Anda untuk data yang tiba setelah window tertutup, dan apakah semua orang di hilir mengetahuinya?
  6. Bagaimana change data capture akan mengubah cara Anda menjaga pencarian, cache, dan analitik tetap tersinkron?

Poin-poin utama

  • Raih streaming hanya ketika keputusan peka-latensi membayarnya; batch dan micro-batch adalah bawaan yang lebih murah.
  • Hitung berdasarkan waktu peristiwa, dan perlakukan data terlambat dan tak berurutan sebagai masalah inti, ditangani dengan window dan watermark.
  • Pilih pengiriman at-least-once dengan sink idempoten untuk hasil effectively-once daripada exactly-once harfiah di mana-mana.
  • Checkpoint pemrosesan berstatus, batasi status Anda, dan pantau lag konsumen sebagai metrik utama.
  • Pakai change data capture untuk streaming dari basis data operasional alih-alih polling.
  • Pilih arsitektur streaming-first di atas log yang disimpan dan dapat diputar ulang daripada memelihara dua basis kode.
  • Ekspos stream lewat streaming SQL, materialized view, dan OLAP real-time, dan jaga setiap stream dapat diamati dan diaudit.

Referensi dan bacaan lanjutan

  • Tyler Akidau, Slava Chernyak, dan Reuven Lax, “Streaming Systems.”
  • Martin Kleppmann, “Designing Data-Intensive Applications.”
  • Nathan Marz dan James Warren, “Big Data” (arsitektur Lambda).
  • Jay Kreps, “Questioning the Lambda Architecture” (O’Reilly Radar).
  • Fabian Hueske dan Vasiliki Kalavri, “Stream Processing with Apache Flink.”
  • Ben Stopford, “Designing Event-Driven Systems.”
  • Tyler Akidau dan rekan, “The Dataflow Model” (makalah VLDB tentang windowing dan watermark).