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

create_table

備考

ベータ版

この機能はベータ版です。

パイプラインでcreate_table()関数を使用し、1つ以上のappend_flow宣言によって書き込まれるマネージドテーブルを作成します。create_table()呼び出しを、テーブルに書き込む1つ以上の@append_flow(target=...)デコレーターとペアリングします。複数のフローが同じマネージドテーブルをターゲットとすることができます。

SQL の同等物については、CREATE TABLE ... FLOW を参照してください。

構文

Python
from pyspark import pipelines as dp

dp.create_table(
name = "<table-name>",
comment = "<comment>",
spark_conf={&quot;&lt;key&gt;&quot; : &quot;&lt;value&gt;&quot;, &quot;&lt;key&gt;&quot; : &quot;&lt;value&gt;&quot;},
table_properties={&quot;&lt;key&gt;&quot; : &quot;&lt;value&gt;&quot;, &quot;&lt;key&gt;&quot; : &quot;&lt;value&gt;&quot;},
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

説明

name

str

必須。テーブル名。

comment

str

テーブルの説明

spark_conf

dict

このクエリーの実行用Spark構成のリスト。

table_properties

dict

テーブルのテーブルプロパティdictです。

partition_cols

list

テーブルのパーティション分割に使用する1つ以上の列のリスト。

path

str

テーブルデータの保存場所。設定されていない場合は、テーブルを含むスキーマのマネージドストレージの場所を使用します。

schema

str または StructType

テーブルのスキーマ定義です。スキーマは、SQL DDL 文字列として、または Python StructType を使用して定義できます。

expect_allexpect_all_or_dropexpect_all_or_fail

dict

テーブルのデータ品質制約。期待値デコレーター関数と同じ動作を提供し、同じ構文を使用しますが、パラメーターとして実装されます。 エクスペクテーションを参照してください。

cluster_by

list

テーブルでリキッドクラスタリングを有効にし、クラスタリングキーとして使用する列を定義します。テーブルにリキッドクラスタリングを使用するを参照してください。

cluster_by_auto

bool

テーブルで自動リキッドクラスタリングを有効にするcluster_byと組み合わせることで、初期クラスタリングキーを定義できます。自動リキッドクラスタリングを参照してください。

row_filter

str

(パブリックプレビュー)テーブルの行フィルター句。「行フィルターと列マスクを使用してテーブルを公開する」を参照してください。

private

bool

Trueの場合、カタログに公開されず、パイプライン内でのみアクセスできるプライベートテーブルを作成します。defaultはFalseです。

パラメーター

Type

説明

name

str

必須。テーブル名。

comment

str

テーブルの説明

spark_conf

dict

このクエリーの実行用Spark構成のリスト。

table_properties

dict

テーブルのテーブルプロパティdictです。

partition_cols

list

テーブルのパーティション分割に使用する1つ以上の列のリスト。

path

str

テーブルデータの保存場所。設定されていない場合は、テーブルを含むスキーマのマネージドストレージの場所を使用します。

schema

str または StructType

テーブルのスキーマ定義です。スキーマは、SQL DDL 文字列として、または Python StructType を使用して定義できます。

expect_allexpect_all_or_dropexpect_all_or_fail

dict

テーブルのデータ品質制約。期待値デコレーター関数と同じ動作を提供し、同じ構文を使用しますが、パラメーターとして実装されます。 エクスペクテーションを参照してください。

cluster_by

list

テーブルでリキッドクラスタリングを有効にし、クラスタリングキーとして使用する列を定義します。テーブルにリキッドクラスタリングを使用するを参照してください。

cluster_by_auto

bool

テーブルで自動リキッドクラスタリングを有効にするcluster_byと組み合わせることで、初期クラスタリングキーを定義できます。自動リキッドクラスタリングを参照してください。

row_filter

str

(パブリックプレビュー)テーブルの行フィルター句。「行フィルターと列マスクを使用してテーブルを公開する」を参照してください。

private

bool

Trueの場合、カタログに公開されず、パイプライン内でのみアクセスできるプライベートテーブルを作成します。defaultはFalseです。

制限事項

  • マネージドテーブルはチェンジデータキャプチャ(CDC)の変更フローをサポートしていません。マネージドテーブルをターゲットとするcreate_auto_cdc_flow()またはcreate_auto_cdc_from_snapshot_flow()が失敗します。CDCターゲットにはcreate_streaming_table()を使用します。
  • マネージドテーブルは append_flow のみをサポートしています。置換フロー (replace_flow / FLOW ... REPLACE WHERE) はサポートされていません。
  • マネージドテーブルは、Unity Catalogを使用するパイプラインでのみサポートされています。
  • 既存のストリーミングテーブルの名前をマネージドテーブルに再利用することはできません。

Python
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")
このページの見出し