メインコンテンツまでスキップ

replace_flow

備考

ベータ版

この機能はベータ版です。

@dp.replace_flow デコレーターは、パイプライン内のストリーミングテーブルに対して REPLACE USING フローを作成します。更新のたびに、フローはターゲットテーブル内の replace_using キー列と一致するすべての行を置き換え、その他の行は変更しません。この関数は、Apache Sparkストリーミング DataFrame を返す必要があります。REPLACE USING フローによる部分的なスナップショットの置換を参照してください。

ソースが列でキー指定された一連の部分的なスナップショットである場合は、@dp.replace_flowを使用します。代わりに単一のステートメントでターゲットテーブルとフローを定義するには、replace_usingsequence_by@dp.tableに渡します。

構文

Python
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>)

パラメーター

パラメーター

Type

説明

function

function

必須。ユーザー定義のクエリから Apache Sparkストリーミング データフレーム を返す関数。

target

str

必須。フローのターゲットであるストリーミングテーブルの名前。

replace_using

list

必須。置換するターゲット行を識別するためのキー列。少なくとも1つの列を指定してください。キー列を重複させることはできず、各キー列の型はソート可能である必要があります。

sequence_by

str または Column

必須。更新の順序を決定する列。各キーについて、最も高いシーケンスが優先され、ターゲットに既に存在する高いシーケンスの行が、低いシーケンスの行によって上書きされることはありません。

name

str

フロー名。指定されていない場合は、defaultで関数名が使用されます。

comment

str

フローの説明。

spark_conf

dict

このクエリーを実行するための Spark 構成の一覧です。

パラメーター

Type

説明

function

function

必須。ユーザー定義のクエリから Apache Sparkストリーミング データフレーム を返す関数。

target

str

必須。フローのターゲットであるストリーミングテーブルの名前。

replace_using

list

必須。置換するターゲット行を識別するためのキー列。少なくとも1つの列を指定してください。キー列を重複させることはできず、各キー列の型はソート可能である必要があります。

sequence_by

str または Column

必須。更新の順序を決定する列。各キーについて、最も高いシーケンスが優先され、ターゲットに既に存在する高いシーケンスの行が、低いシーケンスの行によって上書きされることはありません。

name

str

フロー名。指定されていない場合は、defaultで関数名が使用されます。

comment

str

フローの説明。

spark_conf

dict

このクエリーを実行するための Spark 構成の一覧です。

Python
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")

レコードが列の組み合わせによって識別される場合は、複数のキー列を使用します。

Python
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")

制限事項

  • ストリーミングテーブルは単一の REPLACE USING フローをサポートしており、REPLACE USING を追加フロー、Auto CDC フロー、または REPLACE WHERE フローなどの他のフロータイプと組み合わせることはできません。
  • クエリーはストリーミングクエリーである必要があります。@dp.replace_flow は、ストリーミングではないソースを拒否します。
  • REPLACE USING フローには Databricks Runtime 18.2 以降が必要です。
このページの見出し