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

スタンドアロンのパイプラインでPythonを使用する

Python を使用して、ノートブックからスタンドアロンのマテリアライズドビューやストリーミングテーブルを作成および更新できます。これにより、他の Python ベースのノートブックワークフローと並行してスタンドアロンパイプラインを管理できるようになります。

これを行うには 2 つの方法があります。

  • pyspark.pipelines デコレータ、@dp.materialized_view、および @dp.table を使用してテーブルを定義します。ロジックを DataFrame コードとして表現する方が簡単な場合は、これを使用します。See Define tables with the pipelines decorators.
  • Databricks SQL warehouseが実行するのと同じ SQL ステートメントを spark.sql() に渡して送信します。これにより、REFRESH ステートメントや更新スケジュールなど、スタンドアロンのマテリアライズドビューおよびストリーミングテーブルの SQL サーフェス全体を利用できるようになります。spark.sql() を使用した SQL ステートメントの送信を参照してください。

スタンドアロン パイプライン用の Python ソースには、 サーバレス汎用コンピュート にアタッチされたノートブックが必要です。Databricks SQLウェアハウスから、Pythonを使用してスタンドアロンのパイプラインを作成または更新することはできません。ウェアハウスはPythonノートブックではなくSQLステートメントを実行するためです。代わりにSQLウェアハウスを使用するには、「スタンドアロンのマテリアライズドビューを使用する」および「スタンドアロンのストリーミングテーブルを使用する」を参照してください。

備考

ベータ版

スタンドアロンのマテリアライズドビューとストリーミングテーブルをノートブックからサーバレス汎用コンピュート上で作成および更新することは、ベータ版として、一部の地域で利用可能です。See ノートブック.

要件​

Pythonでスタンドアロンのパイプラインを作成および更新するには、Databricks Runtime 18.1以降のサーバレス汎用コンピュートにアタッチされたノートブックが必要です。地域ごとの提供状況と権限を含む、完全な要件のリストについては、「ノートブック」を参照してください。

パイプラインデコレータを使用したテーブルの定義​

LakeFlow Pipelines で使用するのと同じデコレータを使用して、スタンドアロンのマテリアライズドビューまたはストリーミングテーブルを定義できます。デコレータが適用された各関数は、1 つのテーブルを定義します。セルを実行すると、Databricks はテーブルを作成し、それを入力するためのServerlessパイプラインを実行します。更新が完了すると、セルは結果を返します。

警告

パイプラインデコレーターを使用するには、Serverless 環境のバージョン 5 以上が必要です。

マテリアライズドビューの定義​

バッチ DataFrame を返す関数で @dp.materialized_view を使用します。Wanderbricks サンプルデータセットの bookings テーブルからマテリアライズドビュー daily_booking_revenue を作成する次の例:

Python
from pyspark import pipelines as dp
from pyspark.sql import functions as F

@dp.materialized_view(name="main.default.daily_booking_revenue")
def daily_booking_revenue():
return (
spark.read.table("samples.wanderbricks.bookings")
.groupBy("check_in")
.agg(F.sum("total_amount").alias("total_revenue"))
)

ストリーミング読み取りからテーブルを定義するには、代わりに @dp.table を使用します。

ストリーミングテーブルの定義​

ストリーミング DataFrame を返す関数で @dp.table を使用します。次の例では、同じ bookings テーブルのストリーミング読み取りから、ストリーミングテーブル bookings_raw を作成します。

Python
from pyspark import pipelines as dp

@dp.table(name="main.default.bookings_raw")
def bookings_raw():
return spark.readStream.table("samples.wanderbricks.bookings").select("booking_id", "property_id", "total_amount")

関数がバッチ DataFrame を返す場合、@dp.table は代わりにマテリアライズドビューを作成します。唯一の例外は replace_where であり、常にストリーミングテーブルになります。次の例では、2025 年 7 月 1 日以降のチェックインに関する日次収益を、過去の日付を再計算せずに最新の状態に保ちます。

Python
from pyspark import pipelines as dp
from pyspark.sql import functions as F

@dp.table(
name="main.default.booking_revenue_rw",
replace_where=F.col("check_in") >= F.to_date(F.lit("2025-07-01")),
)
def booking_revenue_rw():
return (
spark.read.table("samples.wanderbricks.bookings")
.groupBy("check_in")
.agg(F.sum("total_amount").alias("total_revenue"))
)

各ランは、述語に一致する行を削除し、その範囲のみを再計算します。REPLACE WHERE フローでのバッチ処理を参照してください。

テーブルの更新​

デコレータを使用して定義したテーブルを更新するには、定義するコードを再度実行します。たとえば、ノートブックのセルを再実行したり、ノートブック全体を実行したり、ノートブックをジョブとして実行したりします。各実行により、テーブルが存在しない場合は作成され、存在する場合は更新されます。

ソースで使用可能なすべてのデータを再処理するには、いずれかのデコレーターに full_refresh=True を渡します。

Python
@dp.table(name="main.default.bookings_raw", full_refresh=True)
def bookings_raw():
return spark.readStream.table("samples.wanderbricks.bookings").select("booking_id", "property_id", "total_amount")

デコレータで定義されたテーブルに対して REFRESH ステートメントを使用したり、SCHEDULE や TRIGGER ON UPDATE で更新をスケジュールしたりすることはできません。スケジュールに基づいて更新するには、SQL でテーブルを定義するか、ノートブックをジョブとしてスケジュールします。See Lakeflow Jobs.

テーブルの構成​

デコレータは、comment、table_properties、partition_cols、cluster_by、schema、spark_conf を含め、パイプライン内部で受け入れるものと同じ共通のデータセットパラメータを受け入れます:

Python
@dp.materialized_view(
name="main.default.daily_booking_revenue",
comment="Daily booking revenue.",
table_properties={"quality": "gold"},
cluster_by=["check_in"],
)
def daily_booking_revenue():
return (
spark.read.table("samples.wanderbricks.bookings")
.groupBy("check_in")
.agg(F.sum("total_amount").alias("total_revenue"))
)

パラメータリストについては、materialized_view および table を参照してください。

private=True プライベートテーブルは同じパイプライン内の他のデータセットからしか読み取ることができないため、サポートされていません。

サポートされていない APIs​

スタンドアロンテーブルは単一のフローを持つ単一のデータセットであるため、データセット間の関係を記述する APIs は使用できません。パイプラインの外部では、次の項目でエラーが発生します:

  • @dp.temporary_view そして dp.create_streaming_table
  • @dp.append_flow およびその他の追加のフロー
  • dp.create_auto_cdc_flow そして dp.create_auto_cdc_from_snapshot_flow
  • @dp.replace_flow および REPLACE USING フローを定義する replace_using パラメーター。REPLACE USING フローによる部分的なスナップショットの置き換えを参照してください。
  • dp.create_sink
  • @dp.expect などのエクスペクテーションと @dp.expect_or_fail

これらを使用するには、代わりに Lakeflow パイプラインを作成します。「Python を使用したパイプライン コードの開発」を参照してください。

次を使用してSQLステートメントを送信: spark.sql()​

Pythonノートブックでは、Databricks SQLウェアハウスから実行するのと同じステートメントを、spark.sql()に渡します。スタンドアロンのマテリアライズドビューとストリーミングテーブルの構文は同一です;ステートメントの送信方法のみが異なります。ウェアハウスと同様に、各CREATEまたはREFRESHステートメントは、操作を処理するためにサーバレス パイプラインを実行します。

spark セッションは Databricks ノートブックでデフォルトで利用できるため、インポートは必要ありません。

マテリアライズドビューの作成​

次の例では、ベーステーブル base_table1 から mv1 マテリアライズドビューを作成します:

Python
spark.sql("""
CREATE OR REPLACE MATERIALIZED VIEW mv1
AS SELECT
date,
sum(sales) AS sum_of_sales
FROM base_table1
GROUP BY date
""")

スケジュールされた更新やトリガーされた更新など、完全なCREATE MATERIALIZED VIEW詳細については、「マテリアライズドビューの作成」を参照してください。

ストリーミングテーブルを作成​

次の例では、raw_dataテーブルからsalesというストリーミングテーブルを作成します:

Python
spark.sql("""
CREATE OR REFRESH STREAMING TABLE sales
AS SELECT product, price FROM STREAM raw_data
""")

CREATE STREAMING TABLE の詳細 (Auto Loader を使用したファイルの読み込みやスケジューリングなど) については、「スタンドアロン ストリーミングテーブルの使用」を参照してください。

マテリアライズドビューまたはストリーミングテーブルを更新​

スタンドアロンテーブルを、ソースからの最新データで更新するには、REFRESHステートメントを使用します。

Python
spark.sql("REFRESH MATERIALIZED VIEW mv1")
spark.sql("REFRESH STREAMING TABLE sales")

サーバレス汎用コンピュートでは、更新は同期的に行われます。非同期更新 (ASYNC キーワード) はサポートされていません。サーバレスジェネラルコンピュートを参照してください。

ステートメントのパラメータ化​

値をPythonコードからステートメントにハードコーディングする代わりに渡すには、SQLで名前付きパラメーターマーカーを使用し、spark.sql()のargs引数を通じてその値を供給します。リテラル値には、:min_salesなどのマーカーを直接使用してください。識別子はプレーンな文字列値として置換できないため、パラメーターがテーブル、ビュー、スキーマなどのオブジェクト名である場合にのみ、マーカーをIDENTIFIER()で囲んでください。

次の例では、マテリアライズドビュー名とフィルター値の両方をパラメーター化します。

Python
mv_name = "main.sales.regional_sales"
min_sales = 1000

spark.sql("""
CREATE OR REPLACE MATERIALIZED VIEW IDENTIFIER(:mv)
AS SELECT
region,
sum(sales) AS sum_of_sales
FROM base_table1
WHERE sales > :min_sales
GROUP BY region
""", args={
"mv": mv_name,
"min_sales": min_sales,
})

詳細については、「パラメーター・マーカー」および「IDENTIFIER 句」を参照してください。

その他のステートメントを実行​

Python ノートブックから、spark.sql() に渡すことで、更新のスケジュール設定、テーブルの変更、またはテーブルの削除などのステートメントを含む、任意のスタンドアロンのマテリアライズドビューまたはストリーミングテーブルのステートメントを実行できます。マテリアライズドビューとストリーミングテーブルの使用方法をSQL構文を含めて理解するには、「スタンドアロンのマテリアライズドビューを使用する」および「スタンドアロンのストリーミングテーブルを使用する」を参照してください。

制限事項:​

サーバレス汎用コンピュートで作成されたスタンドアロンのマテリアライズドビューとストリーミングテーブルには、非同期更新のサポートがないこと、およびテーブルごとのコスト配分がないことなど、追加の制限があります。詳細については、「サーバレス全般コンピュート」を参照してください。

これらのパイプラインはSQL WarehouseではなくServerless汎用コンピュート上でランされるため、囲んでいるwarehouseからカスタムタグを継承しません。system.billing.usage へのwarehouse タグの伝播は、SQL Warehouseからステートメントが実行されるマテリアライズドビューおよびストリーミングテーブルにのみ適用されます。カスタムタグを使用してSQLウェアハウスにコストを配賦するを参照してください。

パイプラインデコレーターで定義されたテーブルには、次の追加の制限があります。

  • REFRESH ステートメントで更新したり、SCHEDULE または TRIGGER ON UPDATE で更新をスケジュールしたりすることはできません。テーブルの更新を参照してください。
  • 期待値、追加フロー、チェンジデータキャプチャ(CDC)フロー、シンク、および一時ビューはサポートされていません。サポートされていない APIs を参照してください。
  • private=True はサポートされていません。

その他のリソース​