replace_flow
ベータ版
この機能はベータ版です。
@dp.replace_flow デコレーターは、パイプライン内のストリーミングテーブルに対して REPLACE USING フローを作成します。更新のたびに、フローはターゲットテーブル内の replace_using キー列と一致するすべての行を置き換え、その他の行は変更しません。この関数は、Apache Sparkストリーミング DataFrame を返す必要があります。REPLACE USING フローによる部分的なスナップショットの置換を参照してください。
ソースが列でキー指定された一連の部分的なスナップショットである場合は、@dp.replace_flowを使用します。代わりに単一のステートメントでターゲットテーブルとフローを定義するには、replace_usingとsequence_byを@dp.tableに渡します。
構文
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 |
| 必須。ユーザー定義のクエリから Apache Sparkストリーミング データフレーム を返す関数。 |
|
| 必須。フローのターゲットであるストリーミングテーブルの名前。 |
|
| 必須。置換するターゲット行を識別するためのキー列。少なくとも1つの列を指定してください。キー列を重複させることはできず、各キー列の型はソート可能である必要があります。 |
|
| 必須。更新の順序を決定する列。各キーについて、最も高いシーケンスが優先され、ターゲットに既に存在する高いシーケンスの行が、低いシーケンスの行によって上書きされることはありません。 |
|
| フロー名。指定されていない場合は、defaultで関数名が使用されます。 |
|
| フローの説明。 |
|
| このクエリーを実行するための Spark 構成の一覧です。 |
例
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")
レコードが列の組み合わせによって識別される場合は、複数のキー列を使用します。
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 以降が必要です。