replace_flow

Important

この機能は ベータ版です

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

@dp.replace_flow元が列ごとにキー化された部分的なスナップショットの連続であれば使います。 ターゲットテーブルとフローを1つの文で定義するには、 replace_usingsequence_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以降が必要です。