Penggantian sebagian cuplikan dengan alur kerja REPLACE USING

Important

Fitur ini ada di Beta.

Alur REPLACE using menjaga tabel target tetap sinkron dengan sumber streaming: menggantikan semua baris yang sesuai dengan kolom kunci yang ditentukan dan membiarkan semua data lain tetap sama.

SEQUENCE BY kolom mengurutkan pembaruan sehingga hasilnya tetap benar meskipun pembaruan diterima dalam urutan yang tidak sesuai. Untuk setiap kunci, urutan tertinggi menang, dan baris dengan urutan lebih rendah tidak pernah menimpa baris yang lebih tinggi yang sudah ada di target. Baris yang berbagi kunci dan urutan yang sama akan ditambahkan, bukan diganti.

Cara kerja REPLACE USING

Pertimbangkan sebuah tabel peristiwa yang memuat event klik dan konversi untuk dua wilayah, yang diurutkan oleh seq:

region_id jenis perangkat event_type seq
1 iOS klik 1
1 Android konversi 1
2 iOS klik 1
2 desktop klik 1

Sebuah REPLACE USING (region_id) SEQUENCE BY seq aliran menerima pembaruan ini untuk wilayah 1 dan 3. Region 2 tidak memiliki pembaruan:

region_id jenis perangkat event_type seq
1 iOS klik 2
1 Android konversi 2
1 desktop klik 2
3 iOS klik 1
3 desktop klik 2

Targetnya menjadi:

region_id jenis perangkat event_type seq Hasil
1 iOS klik 2 Diganti, karena seq 2 lebih besar dari seq 1
1 Android konversi 2 Diganti, karena seq 2 lebih besar dari seq 1
1 desktop klik 2 Diganti, karena seq 2 lebih besar dari seq 1
2 iOS klik 1 Tidak tersentuh, karena kunci tidak ada dalam pembaruan ini
2 desktop klik 1 Tidak tersentuh, karena kunci tidak ada dalam pembaruan ini
3 desktop klik 2 Ditambahkan. Baris seq 1 untuk wilayah 3 tidak ditambahkan, karena hanya urutan tertinggi untuk suatu kunci yang diterapkan.

Persyaratan

Alur REPLACE USING memiliki persyaratan berikut:

  • REPLACE USING alur dijalankan di Databricks Runtime 18.2 dan yang lebih baru, pada komputasi klasik atau serverless. Databricks merekomendasikan Unity Catalog.
  • Sumber harus berupa sumber streaming. REPLACE USING menolak sumber non-streaming.
  • Anda harus menentukan setidaknya satu kolom kunci dan tepat satu kolom SEQUENCE BY.

Kapan menggunakan alur REPLACE USING

Pipa Lakeflow menawarkan tiga aliran yang menimpa baris yang sudah ada. Pilih berdasarkan seperti apa sumber Anda dan bagaimana ia mengidentifikasi baris yang akan diganti:

  • Gunakan REPLACE USING saat sumber Anda berupa serangkaian snapshot parsial yang diidentifikasi berdasarkan kolom. REPLACE USING hanya menimpa data yang memiliki kecocokan dalam data masuk, sehingga semua data lain tetap tidak tersentuh. Tidak memerlukan kunci utama.
  • Gunakan AUTO CDC ketika sumber Anda adalah feed penangkapan data perubahan (CDC) dengan operasi sisipkan, pembaruan, dan hapus eksplisit, atau Anda membutuhkan riwayat Tipe 2 dimensi yang berubah perlahan (SCD). AUTO CDC juga membutuhkan kunci utama yang sesungguhnya. Lihat API CDC Otomatis: Menyederhanakan penangkapan perubahan data menggunakan pipeline.
  • Gunakan REPLACE WHERE jika sumber Anda adalah snapshot dan Anda ingin menghitung kembali serta menimpa rentang pada tabel target yang dipilih berdasarkan predikat, misalnya 7 hari terakhir, dalam operasi batch. Tidak memerlukan kunci utama. Lihat Pemrosesan batch dengan alur REPLACE WHERE.

Buat alur GANTI MENGGUNAKAN

Definisikan alur REPLACE USING baik di SQL maupun Python.

SQL

Gunakan klausa sebaris FLOW REPLACE USING dengan CREATE STREAMING TABLE:

CREATE STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

Atau, gunakan sintaks bentuk CREATE FLOW panjang:

CREATE STREAMING TABLE payments_current;

CREATE FLOW payments_flow AS
INSERT INTO payments_current BY NAME
REPLACE USING (payment_id) SEQUENCE BY payment_date
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

Note

BY NAME diperlukan dalam SQL. Ini mencocokkan kolom berdasarkan nama, bukan berdasarkan posisi.

Python

Deklarasikan tabel dan alur bersama dengan @dp.table:

from pyspark import pipelines as dp

@dp.table(name="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_current():
  return spark.readStream.table("samples.wanderbricks.payments")

Sebagai alternatif, targetkan tabel streaming yang sudah ada dengan @dp.replace_flow:

from pyspark import pipelines as dp

dp.create_streaming_table("payments_current")

@dp.replace_flow(target="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_flow():
  return spark.readStream.table("samples.wanderbricks.payments")

replace_using adalah daftar kolom kunci. sequence_by adalah nama kolom atau ekspresi Column, dan diperlukan setiap kali replace_using ditetapkan.

Pengurutan dan data yang tidak berurutan

Kolom SEQUENCE BY membuat hasilnya tidak bergantung pada urutan pembaruan diterima. Sebuah baris diterapkan pada kunci hanya jika urutannya lebih besar dari urutan yang sudah disimpan untuk kunci tersebut, sehingga baris yang terlambat atau diputar ulang yang lebih tua dari nilai saat ini diabaikan. Tombol yang tidak ada dalam pembaruan tidak akan disentuh.

Ikuti praktik-praktik ini agar penggantian berperilaku dapat diprediksi:

Praktik Alasan
Gunakan urutan yang selalu bertambah untuk setiap versi kunci, seperti stempel waktu, nomor versi, atau offset log. Dua baris dengan kunci dan urutan yang sama keduanya dipertahankan, yang menghasilkan baris duplikat untuk kunci tersebut.
Gunakan urutan non-null. Urutan null dapat menyebabkan perilaku yang tidak terdefinisi.

Expectations

GANTI MENGGUNAKAN alur mendukung ekspektasi. warn dan fail berperilaku seperti pada alur lain: warn terus melanggar baris dan merekam pelanggaran, serta fail menghentikan pembaruan. Lihat Mengelola kualitas data dengan ekspektasi alur kerja.

Sebuah ekspektasi drop memperlakukan baris yang melanggar seolah-olah sumber tidak pernah menghasilkannya. Baris yang dihapus tidak menggantikan, menghapus, atau memodifikasi kunci yang cocok di tabel target:

  • Penghapusan dilakukan sebelum deduplikasi, sehingga alur mempertahankan versi valid terbaru untuk kunci tersebut.
  • Jika setiap baris masuk untuk sebuah kunci dihapus, baris kunci yang ada akan tetap tidak tersentuh.
  • Karena baris yang dihapus tidak menetapkan lantai urutan, pembaruan valid yang lebih baru tetap muncul meskipun urutannya lebih rendah dari baris yang dihapus.

Keterbatasan

Alur GANTI MENGGUNAKAN memiliki keterbatasan berikut:

  • REPLACE USING mendukung satu alur untuk setiap tabel target. Menggabungkan REPLACE USING dengan tipe aliran lain pada target yang sama tidak didukung.
  • Tabel target harus dibuat di dalam pipeline.
  • Sumber harus berupa sumber streaming.
  • Anda harus menentukan setidaknya satu kolom kunci dan satu SEQUENCE BY kolom. Kolom kunci tidak dapat diulang, dan tipe setiap kolom kunci harus dapat diurutkan. Tipe atom, seperti bilangan bulat, string, dan tanggal, dapat berupa kunci, sementara MAP dan VARIANT tidak bisa.
  • Untuk tabel streaming mandiri, lihat Menerapkan penggantian snapshot parsial dengan flow REPLACE USING untuk mengetahui perbedaan sintaks.

Examples

Contoh berikut dibaca dari samples.wanderbricks.booking_updates, sebuah tabel contoh perubahan status pemesanan yang tersedia di setiap workspace yang mendukung Unity Catalog. Setiap pemesanan muncul satu kali untuk setiap perubahan, sehingga booking_id muncul kembali dengan booking_update_id yang baru. Lihat dataset Wanderbricks.

Contoh 1: Simpan catatan terbaru untuk setiap kunci

Pertahankan hanya status terkini dari setiap pemesanan. Flow ini menggunakan booking_id sebagai kunci dan mengurutkan berdasarkan booking_update_id, sehingga pembaruan terbaru untuk pemesanan menggantikan pembaruan yang lebih lama. Gunakan AUTO CDC sebagai gantinya saat sumber Anda berupa aliran perubahan dengan operasi penyisipan, pembaruan, dan penghapusan yang eksplisit.

SQL

CREATE OR REFRESH STREAMING TABLE bookings_current
FLOW REPLACE USING (booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_current",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
def bookings_current():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

Contoh ini diurutkan berdasarkan booking_update_id alih-alih stempel waktu updated_at, karena beberapa pembaruan pada pemesanan yang sama dapat memiliki stempel waktu yang sama. Baris yang berhubungan pada urutan akan ditambahkan, bukan diganti, sehingga akan menyisakan lebih dari satu baris untuk pemesanan tersebut.

Contoh 2: Kunci pada lebih dari satu kolom

Ketika sebuah catatan diidentifikasi dengan kombinasi kolom, daftarkan semuanya dalam .REPLACE USING Di sini setiap pemesanan diidentifikasi dengan (property_id, booking_id), sehingga alur tetap mempertahankan status terkini dari setiap pemesanan per properti. Jika kolom kunci dapat bernilai null, REPLACE USING mencocokkan null dengan null alih-alih melewati baris tersebut.

SQL

CREATE OR REFRESH STREAMING TABLE bookings_by_property
FLOW REPLACE USING (property_id, booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT property_id, booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_by_property",
  replace_using=["property_id", "booking_id"],
  sequence_by="booking_update_id"
)
def bookings_by_property():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

Contoh 3: Menghapus catatan tidak valid dengan ekspektasi

Tambahkan ekspektasi untuk menjaga baris yang buruk tetap di luar target. Baris yang dihilangkan diperlakukan seolah-olah sumber tidak pernah memproduksinya: tidak menggantikan atau menghapus kunci yang cocok, dan alur akan kembali ke baris valid terbaru untuk kunci tersebut. Alur ini akan menghapus pembaruan yang tidak memiliki nilai positif total_amount.

from pyspark import pipelines as dp

@dp.table(
  name="bookings_validated",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
@dp.expect_or_drop("positive_amount", "total_amount > 0")
def bookings_validated():
  return spark.readStream.table("samples.wanderbricks.booking_updates")