Catatan
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba masuk atau mengubah direktori.
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba mengubah direktori.
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.
Sebuah kolom mengatur SEQUENCE BY pembaruan sehingga hasilnya benar meskipun pembaruan datang tidak sesuai urutan. 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. Seq 1 baris untuk region 3 tidak ditambahkan, karena hanya urutan tertinggi untuk sebuah kunci yang diterapkan. |
Persyaratan
GANTI MENGGUNAKAN alur memiliki persyaratan berikut:
- GANTI MENGGUNAKAN alur yang dijalankan pada Databricks Runtime 18.2 ke atas, pada komputasi klasik atau serverless. Databricks merekomendasikan Unity Catalog.
- Sumber harus berupa sumber streaming. GANTI MENGGUNAKAN menolak sumber non-streaming.
- Anda harus menentukan setidaknya satu kolom kunci dan tepat satu
SEQUENCE BYkolom.
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 ketika sumber Anda adalah serangkaian snapshot parsial yang dikunci 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 ketika sumber Anda adalah snapshot dan Anda ingin menghitung ulang serta menimpa rentang tabel target yang dipilih oleh predikat, misalnya 7 hari terakhir, sebagai 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 Column ekspresi, dan diperlukan setiap kali replace_using diatur.
Sekuensing dan data di luar urutan
Kolom membuat SEQUENCE BY hasil tersebut tidak bergantung pada urutan pembaruan yang tiba. 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 meningkat secara ketat per versi kunci, seperti timestamp, 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 drop ekspektasi memperlakukan pertengkaran yang melanggar seolah-olah sumber tidak pernah membawanya. Baris yang dihapus tidak menggantikan, menghapus, atau memodifikasi kunci yang cocok di tabel target:
- Dropping terjadi sebelum deduplikasi, jadi alur tetap menyimpan versi valid terbaru untuk kunci.
- 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
GANTI MENGGUNAKAN alur memiliki keterbatasan berikut:
- REPLACE USING mendukung satu tabel aliran per 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 BYkolom. Kolom kunci tidak dapat diulang, dan tipe setiap kolom kunci harus dapat diurutkan. Tipe atom, seperti bilangan bulat, string, dan tanggal, dapat berupa kunci, sementaraMAPdanVARIANTtidak bisa.
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 sekali setiap perubahan, sehingga booking_id diulang dengan perubahan baru booking_update_id. Lihat dataset Wanderbricks.
Contoh 1: Simpan catatan terbaru untuk setiap kunci
Pertahankan hanya status terkini dari setiap pemesanan. Tombol flow aktif booking_id dan urutan dengan booking_update_id, sehingga pembaruan terbaru untuk pemesanan menggantikan pembaruan sebelumnya. Gunakan AUTO CDC sebagai gantinya ketika sumber Anda adalah feed perubahan dengan operasi sisipkan, pembaruan, dan penghapusan 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 berurutan berdasarkan booking_update_idupdated_at timestamp karena beberapa pembaruan pada pemesanan yang sama dapat berbagi timestamp. 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 bisa bernilai null, GANTI MENGGUNAKAN mencocokkan null dengan null daripada melewati baris.
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")