create_table
ベータ版
この機能はベータ版です。
パイプラインでcreate_table()関数を使用し、1つ以上のappend_flow宣言によって書き込まれるマネージドテーブルを作成します。create_table()呼び出しを、テーブルに書き込む1つ以上の@append_flow(target=...)デコレーターとペアリングします。複数のフローが同じマネージドテーブルをターゲットとすることができます。
SQL の同等物については、CREATE TABLE ... FLOW を参照してください。
構文
from pyspark import pipelines as dp
dp.create_table(
name = "<table-name>",
comment = "<comment>",
spark_conf={"<key>" : "<value>", "<key>" : "<value>"},
table_properties={"<key>" : "<value>", "<key>" : "<value>"},
partition_cols=["<partition-column>", "<partition-column>"],
path="<storage-location-path>",
schema="schema-definition",
expect_all = {"<key>" : "<value>", "<key>" : "<value>"},
expect_all_or_drop = {"<key>" : "<value>", "<key>" : "<value>"},
expect_all_or_fail = {"<key>" : "<value>", "<key>" : "<value>"},
cluster_by = ["<clustering-column>", "<clustering-column>"],
cluster_by_auto = False,
row_filter = "row-filter-clause",
private = False
)
パラメーター
パラメーター | Type | 説明 |
|---|---|---|
|
| 必須。テーブル名。 |
|
| テーブルの説明 |
|
| このクエリーの実行用Spark構成のリスト。 |
|
| テーブルのテーブルプロパティの |
|
| テーブルのパーティション分割に使用する1つ以上の列のリスト。 |
|
| テーブルデータの保存場所。設定されていない場合は、テーブルを含むスキーマのマネージドストレージの場所を使用します。 |
|
| テーブルのスキーマ定義です。スキーマは、SQL DDL 文字列として、または Python |
|
| テーブルのデータ品質制約。期待値デコレーター関数と同じ動作を提供し、同じ構文を使用しますが、パラメーターとして実装されます。 エクスペクテーションを参照してください。 |
|
| テーブルでリキッドクラスタリングを有効にし、クラスタリングキーとして使用する列を定義します。テーブルにリキッドクラスタリングを使用するを参照してください。 |
|
| テーブルで自動リキッドクラスタリングを有効にする |
|
| (パブリックプレビュー)テーブルの行フィルター句。「行フィルターと列マスクを使用してテーブルを公開する」を参照してください。 |
|
|
|
制限事項
- マネージドテーブルはチェンジデータキャプチャ(CDC)の変更フローをサポートしていません。マネージドテーブルをターゲットとする
create_auto_cdc_flow()またはcreate_auto_cdc_from_snapshot_flow()が失敗します。CDCターゲットにはcreate_streaming_table()を使用します。 - マネージドテーブルは
append_flowのみをサポートしています。置換フロー (replace_flow/FLOW ... REPLACE WHERE) はサポートされていません。 - マネージドテーブルは、Unity Catalogを使用するパイプラインでのみサポートされています。
- 既存のストリーミングテーブルの名前をマネージドテーブルに再利用することはできません。
例
from pyspark import pipelines as dp
dp.create_table("combined")
@dp.append_flow(target="combined")
def from_a():
return spark.readStream.table("source_a")
@dp.append_flow(target="combined")
def from_b():
return spark.readStream.table("source_b")