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.
Dekorator @dp.replace_flow membuat alur GANTI MENGGUNAKAN untuk tabel streaming di pipeline Anda. Pada setiap pembaruan, alur menggantikan semua baris dalam tabel target yang sesuai dengan replace_using kolom kunci dan membiarkan semua baris lainnya tidak tersentuh. Fungsi harus mengembalikan DataFrame streaming Apache Spark. Lihat Penggantian snapshot parsial dengan alur GANTI MENGGUNAKAN.
Gunakan @dp.replace_flow ketika sumber Anda adalah serangkaian snapshot parsial yang dikunci berdasarkan kolom. Untuk mendefinisikan tabel target dan alur dalam satu pernyataan, replace_using kirim dan sequence_by ke @dp.table.
Syntax
from pyspark import pipelines as dp
dp.create_streaming_table("<target-table-name>") # Required only if the target table doesn't exist.
@dp.replace_flow(
target = "<target-table-name>",
replace_using = ["<key-column>", "<key-column>"],
sequence_by = "<sequence-column>",
name = "<flow-name>", # optional, defaults to function name
comment = "<comment>", # optional
spark_conf = {"<key>" : "<value>", "<key>" : "<value>"}) # optional
def <function-name>():
return (<streaming-query>)
Parameters
| Parameter | Type | Deskripsi |
|---|---|---|
| fungsi | function |
Required. Fungsi yang mengembalikan DataFrame streaming Apache Spark dari kueri yang ditentukan pengguna. |
target |
str |
Required. Nama tabel streaming yang menjadi target aliran. |
replace_using |
list |
Required. Kolom kunci yang mengidentifikasi baris target mana yang akan diganti. Tentukan setidaknya satu kolom. Kolom kunci tidak dapat diulang, dan tipe setiap kolom kunci harus dapat diurutkan. |
sequence_by |
str atau Column |
Required. Kolom yang mengatur pembaruan. Untuk setiap kunci, urutan tertinggi menang, dan baris dengan urutan lebih rendah tidak pernah menimpa baris yang lebih tinggi yang sudah ada di target. |
name |
str |
Nama alur. Jika tidak disediakan, akan otomatis menggunakan nama fungsi. |
comment |
str |
Deskripsi untuk alur. |
spark_conf |
dict |
Daftar konfigurasi Spark untuk eksekusi kueri ini. |
Examples
from pyspark import pipelines as dp
# Keep the latest row for each order from a stream of partial snapshots
dp.create_streaming_table("orders_current")
@dp.replace_flow(
target = "orders_current",
replace_using = ["order_id"],
sequence_by = "updated_at"
)
def orders_flow():
return spark.readStream.table("order_updates")
Gunakan lebih dari satu kolom kunci ketika sebuah catatan diidentifikasi dengan kombinasi kolom:
from pyspark import pipelines as dp
dp.create_streaming_table("accounts_current")
@dp.replace_flow(
target = "accounts_current",
replace_using = ["region", "account_id"],
sequence_by = "updated_at"
)
def accounts_flow():
return spark.readStream.table("account_updates")
Keterbatasan
- Tabel streaming mendukung satu
REPLACE USINGaliran, dan tidak dapat digabungkanREPLACE USINGdengan tipe alur lain seperti alur tambahan, alur CDC otomatis, atau aliranREPLACE WHERE. - Kueri harus berupa kueri streaming.
@dp.replace_flowmenolak sumber non-streaming. -
REPLACE USINGalur memerlukan Databricks Runtime 18.2 dan ke atasnya.