Skip to main content

create_table

Beta

This feature is in Beta.

Use the create_table() function in a pipeline to create a managed table, written by one or more append_flow declarations. Pair the create_table() call with one or more @append_flow(target=...) decorators that write into the table. Multiple flows can target the same managed table.

For the SQL equivalent, see CREATE TABLE ... FLOW.

Syntax​

Python
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
)

Parameters​

Parameter

Type

Description

name

str

Required. The table name.

comment

str

A description for the table.

spark_conf

dict

A list of Spark configurations for the execution of this query.

table_properties

dict

A dict of table properties for the table.

partition_cols

list

A list of one or more columns to use for partitioning the table.

path

str

A storage location for table data. If not set, use the managed storage location for the schema containing the table.

schema

str or StructType

A schema definition for the table. Schemas can be defined as a SQL DDL string or with a Python StructType.

expect_all, expect_all_or_drop, expect_all_or_fail

dict

Data quality constraints for the table. Provides the same behavior and uses the same syntax as expectation decorator functions, but implemented as a parameter. See Expectations.

cluster_by

list

Enable liquid clustering on the table and define the columns to use as clustering keys. See Use liquid clustering for tables.

cluster_by_auto

bool

Enable automatic liquid clustering on the table. Can be combined with cluster_by to define the initial clustering keys. See Automatic liquid clustering.

row_filter

str

(Public Preview) A row filter clause for the table. See Publish tables with row filters and column masks.

private

bool

When True, creates a private table that is not published to the catalog and is only accessible within the pipeline. Defaults to False.

Parameter

Type

Description

name

str

Required. The table name.

comment

str

A description for the table.

spark_conf

dict

A list of Spark configurations for the execution of this query.

table_properties

dict

A dict of table properties for the table.

partition_cols

list

A list of one or more columns to use for partitioning the table.

path

str

A storage location for table data. If not set, use the managed storage location for the schema containing the table.

schema

str or StructType

A schema definition for the table. Schemas can be defined as a SQL DDL string or with a Python StructType.

expect_all, expect_all_or_drop, expect_all_or_fail

dict

Data quality constraints for the table. Provides the same behavior and uses the same syntax as expectation decorator functions, but implemented as a parameter. See Expectations.

cluster_by

list

Enable liquid clustering on the table and define the columns to use as clustering keys. See Use liquid clustering for tables.

cluster_by_auto

bool

Enable automatic liquid clustering on the table. Can be combined with cluster_by to define the initial clustering keys. See Automatic liquid clustering.

row_filter

str

(Public Preview) A row filter clause for the table. See Publish tables with row filters and column masks.

private

bool

When True, creates a private table that is not published to the catalog and is only accessible within the pipeline. Defaults to False.

Limitations​

  • Managed tables do not support change data capture (CDC) change flows. create_auto_cdc_flow() or create_auto_cdc_from_snapshot_flow() targeting a managed table fails. Use create_streaming_table() for CDC targets.
  • Managed tables support only append_flow. Replace flows (replace_flow / FLOW ... REPLACE WHERE) are not supported.
  • Managed tables are supported only in pipelines with Unity Catalog.
  • You cannot reuse the name of an existing streaming table for a managed table.
  • create_table cannot adopt an existing managed table. The target must be a new managed table.

Example​

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")