replace_flow

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 USING aliran, dan tidak dapat digabungkan REPLACE USING dengan tipe alur lain seperti alur tambahan, alur CDC otomatis, atau aliran REPLACE WHERE .
  • Kueri harus berupa kueri streaming. @dp.replace_flow menolak sumber non-streaming.
  • REPLACE USING alur memerlukan Databricks Runtime 18.2 dan ke atasnya.