create_table

Important

此功能在 Beta 版中。

create_table()使用管道中的函数创建由一个或多个append_flow声明编写的托管表。 将 create_table() 调用与写入表中的一个或多个 @append_flow(target=...) 修饰器配对。 多个流可以面向同一托管表。

有关 SQL 等效项,请参阅 CREATE TABLE ...FLOW

Syntax

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

参数 类型 说明
name str Required. 表名称。
comment str 表的说明。
spark_conf dict 用于执行此查询的 Spark 配置列表。
table_properties dict 表的dict
partition_cols list 用于对表进行分区的一个或多个列的列表。
path str 表数据的存储位置。 如果未设置,请使用包含表的架构的托管存储位置。
schema strStructType 表的架构定义。 架构可以定义为 SQL DDL 字符串,或使用 Python StructType 定义。
expect_allexpect_all_or_dropexpect_all_or_fail dict 数据表的质量约束。 提供相同的行为并使用与预期修饰器函数相同的语法,但作为参数实现。 请参阅 期望值。
cluster_by list 对表启用动态聚类,并定义要用作聚类键的列。 请参阅对表使用 liquid 聚类分析
cluster_by_auto bool 启用表格上的自动液体聚类功能。 可以与 cluster_by 定义初始聚类分析键结合使用。 请参阅 自动液体聚类
row_filter str (公共预览版)表的行筛选器子句。 请参阅发布具有行筛选器和列掩码的表
private bool True,创建一个未发布到目录的专用表,并且只能在管道中访问。 默认值为 False.

局限性

  • 托管表不支持变更数据捕获(CDC)更改流。 create_auto_cdc_flow()create_auto_cdc_from_snapshot_flow() 目标托管表失败。 将 create_streaming_table() 用于 CDC 目标。
  • 托管表仅 append_flow支持 。 不支持替换流(replace_flow / FLOW ... REPLACE WHERE)。
  • 托管表仅在具有 Unity 目录的管道中受支持。
  • 不能重复使用托管表的现有流式处理表的名称。

示例

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