Important
この機能は ベータ版です。
@dp.replace_flowデコレーターはパイプライン内のストリーミングテーブル用に「REPLACE USING flow」を作成します。 各更新時には、フローはターゲットテーブル内の replace_using キー列に一致するすべての行を置き換え、他の行はそのままにします。 この関数は、Apache Spark ストリーミング DataFrame を返す必要があります。 部分 スナップショット置換はREPLACE USING flowsで参照してください。
@dp.replace_flow元が列ごとにキー化された部分的なスナップショットの連続であれば使います。 ターゲットテーブルとフローを1つの文で定義するには、 replace_using と sequence_by を @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>)
パラメーター
| パラメーター | タイプ | 説明 |
|---|---|---|
| 関数 | function |
必須。 ユーザー定義クエリから Apache Spark ストリーミング DataFrame を返す関数。 |
target |
str |
必須。 フローの対象となるストリーミングテーブルの名前です。 |
replace_using |
list |
必須。 どのターゲット行を置き換えるべきかを示すキーカラムです。 少なくとも1つの列を指定してください。 キー列は重複できず、各キー列のタイプはソート可能でなければなりません。 |
sequence_by |
str または Column |
必須。 更新を指示する列です。 各キーに対して、最も高い列が勝ち、低い列がターゲット内のより高い列を上書きすることはありません。 |
name |
str |
フロー名。 指定しない場合は、既定で関数名が使用されます。 |
comment |
str |
フローの説明。 |
spark_conf |
dict |
このクエリを実行するための 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フローをサポートしており、追加フロー、自動CDCフロー、REPLACE WHEREフローなど他のフロータイプとREPLACE USINGを組み合わせることはできません。 - クエリはストリーミング クエリである必要があります。
@dp.replace_flowストリーミングでないソースを拒否します。 -
REPLACE USINGフローはDatabricks Runtime 18.2以降が必要です。