append_flow

@dp.append_flowデコレーターは、パイプライン テーブルの追加フローまたはバックフィルを作成します。 この関数は、Apache Spark ストリーミング DataFrame を返す必要があります。 「Lakeflow パイプライン フローを使用してデータを増分的に読み込んで処理する」を参照してください。

アペンドフローはストリーミングテーブル、 管理テーブル、またはシンクを対象にできます。

構文

from pyspark import pipelines as dp

dp.create_streaming_table("<target-table-name>") # Required only if the target table doesn't exist.

@dp.append_flow(
  target = "<target-table-name>",
  name = "<flow-name>", # optional, defaults to function name
  once = False, # optional
  depends_on = "<flow-name>", # optional, Public Preview
  spark_conf = {"<key>" : "<value", "<key" : "<value>"}, # optional
  comment = "<comment>", # optional
  import_checkpoint = "<checkpoint-path>") # optional
def <function-name>():
  return (<streaming-query>) #

パラメーター

パラメーター タイプ Description
関数 function 必須。 ユーザー定義クエリから Apache Spark ストリーミング DataFrame を返す関数。
target str 必須。 追加フローのターゲットであるテーブルまたはシンクの名前。
name str フロー名。 指定しない場合は、既定で関数名が使用されます。
once bool 必要に応じて、バックフィルなどの 1 回限りのフローとしてフローを定義します。 once=Trueを使用すると、次の 2 つの方法でフローが変更されます。
  • 戻り値。 streaming-query。 この場合は、ストリーミング DataFrame ではなく、バッチ DataFrame である必要があります。
  • フローは既定で 1 回実行されます。 パイプラインが完全な更新で更新された場合、 ONCE フローが再度実行され、データが再作成されます。
depends_on str または list パブリック プレビュー。 このフローが開始される前に、1つ以上のフロー名が成功裏に完了する必要があります。 単一のフロー名または名前のリストを受け入れます。 これはフローの実行のみを指示します。フローの実行方法を変えることはありません。 「 depends_onによる注文パイプラインフロー実行」を参照してください。
comment str フローの説明。
spark_conf dict このクエリを実行するための Spark 構成の一覧
import_checkpoint str 既存の構造化ストリーミングチェックポイントへのパスをフローにインポートし、移行したストリームがソースを再処理する代わりに、最後にコミットしたオフセットから再開できるようにします。 チェックポイントのインポートは ベータ版です。 詳細は 「構造化ストリーミングチェックポイントの移行」を参照してください。

例示

from pyspark import pipelines as dp

# Create a sink for an external Delta table
dp.create_sink("my_sink", "delta", {"path": "/tmp/delta_sink"})

# Add an append flow to an external Delta table
@dp.append_flow(name = "flow", target = "my_sink")
def flowFunc():
  return <streaming-query>

# Add a backfill
@dp.append_flow(name = "backfill", target = "my_sink", once = True)
def backfillFlowFunc():
    return (
      spark.read
      .format("json")
      .load("/path/to/backfill/")
    )

# Create a Kafka sink
dp.create_sink(
  "my_kafka_sink",
  "kafka",
  {
    "kafka.bootstrap.servers": "host:port",
    "topic": "my_topic"
  }
)

# Add an append flow to a Kafka sink
@dp.append_flow(name = "flow", target = "my_kafka_sink")
def myFlow():
  return read_stream("xxx").select(F.to_json(F.struct("*")).alias("value"))

構造化ストリーミングチェックポイントの移行

Important

チェックポイントのインポートは ベータ版です。

import_checkpointを使って、既存の構造化ストリーミングワークロードをソースを再処理せずにパイプラインに移行できます。 構造化ストリーミングクエリが使ったもの checkpointLocation 、クラウドストレージ、Unity Catalogボリューム、DBFSパスに設定してください。 最初のパイプライン更新では、そのチェックポイントをパイプラインの管理ストレージにクローンします。 その後、フローは最後にコミットされたオフセットから再開され、その状態(アグリゲーション、重複除去キー、ウォーターマークなど)がそのままです。 その後のパイプライン更新では、フローのクローンチェックポイントを使用します。元のチェックポイントは変更されません。

フローは create_table または シンクで作成された管理されたテーブルをターゲットにしなければなりません。

パイプラインを実行する前に、元のStructured Streamingクエリを停止してください。 元のStructured Streamingクエリはインポート後に再利用できますが、チェックポイントの状態を管理し、パイプラインとクエリが同時に同じテーブルに書き込みないようにする必要があります。重複したデータが発生する可能性があります。

構造化ストリーミングクエリをパイプラインフローとして再作成し、新しいテーブルに書き込み、そのチェックポイントをインポートします:

from pyspark import pipelines as dp

# Create a new managed table for the pipeline
dp.create_table("target_table")

# Continue from the imported checkpoint instead of reprocessing the source.
@dp.append_flow(
  target = "target_table",
  import_checkpoint = "/Volumes/my_catalog/my_schema/checkpoints/my_stream",
)
def migrate():
  # The same source your original query read from.
  return spark.readStream.table("source_table")

チェックポイントは最初のパイプライン更新時に一度だけインポートされます。後の更新では import_checkpointを無視します。 完全なリフレッシュはチェックポイントを再インポートするものではなく、新しい空のチェックポイントからソースを再処理します。 異なるチェックポイントをインポートするには、ターゲットテーブルでこれまで使われていないフロー名を使用します。既存のフロー名を再利用するとインポートが飛ばされます。

制限事項

  • 既存のテーブル(例えば元のStructured Streamingクエリターゲット)にチェックポイントをインポートすることはサポートされていません。 パイプラインが作成する新しいテーブルや シンクをターゲットにしてください。
  • import_checkpoint append_flowのみ対応されています。