CREATE STREAMING TABLE (管道)

流式处理表是一种 Delta 表,它额外支持流式处理或增量数据处理。 流表由管道支撑。 每次刷新流式处理表时,添加到源表的数据都会追加到流式处理表中。 可以手动或按计划刷新流式处理表。

若要详细了解如何执行或计划刷新,请参阅 运行管道更新。

Syntax

CREATE [OR REFRESH] [PRIVATE] STREAMING TABLE
  table_name
  [ table_specification ]
  [ table_clauses ]
  [ {flow_clause | AS query} ]

table_specification
  ( { column_identifier column_type [column_properties] } [, ...]
    [ column_constraint ] [, ...]
    [ , table_constraint ] [...] )

   column_properties
      { NOT NULL | GENERATED ALWAYS AS ( expr ) | GENERATED { ALWAYS | BY DEFAULT } AS IDENTITY [ ( [ START WITH start | INCREMENT BY step ] [ ...] ) ] | DEFAULT default_expression | COMMENT column_comment | column_constraint | MASK clause } [ ... ]

table_clauses
  { USING DELTA
    PARTITIONED BY (col [, ...]) |
    CLUSTER BY clause |
    LOCATION path |
    COMMENT view_comment |
    TBLPROPERTIES clause |
    WITH { ROW FILTER clause } } [ ... ]
   } [ ... ]

flow_clause
  FLOW { { INSERT [ONCE] BY NAME query } |
  { AUTO CDC auto_cdc_flow_spec } |
  { REPLACE WHERE predicate BY NAME query } }

参数

  • 刷新

    如果指定,则创建表或更新现有表及其内容。

  • 专用

    创建专用流式处理表。

    • 它们不会添加到目录中,并且只能在定义管道中访问
    • 它们可以与目录中的现有对象同名。 在管道中,如果专用流式处理表和目录中的对象具有相同的名称,则对名称的引用将解析为专用流式处理表。
    • 专用流式处理表仅在管道的生存期内保留,而不仅仅是单个更新。

    专用流式处理表以前是使用 TEMPORARY 参数创建的。

  • table_name

    新创建的表的名称。 完全限定的表名称必须是独一无二的。

  • table_specification

    此可选子句定义列的列表、列的类型、属性、说明和列约束。

    • column_identifier

      列名必须具有唯一性,并映射到查询的输出列。

    • column_type

      指定列的数据类型。 并非 Azure Databricks 支持的所有数据类型都受流式处理表支持。

    • column_comment

      描述列的可选 STRING 文本。 此选项必须与 column_type 一起指定。 如果未指定列类型,则会跳过列注释。

    • GENERATED ALWAYS AS ( expr )

      指定此子句后,此列的值取决于 expr。

      表的 DEFAULT COLLATION 须为 UTF8_BINARY。

      expr 可能包含文本、表中的列标识符以及内置的确定性 SQL 函数或运算符,但以下内容除外:

      此外,expr 不能包含任何子查询。

    • GENERATED { ALWAYS |默认情况下 } AS IDENTITY [ ( [ START WITH start ] [ INCREMENT BY step ] ] ] ]

      适用于:勾选“是” Databricks SQL 勾选“是” Databricks Runtime 10.4 LTS 及更高版本

      定义标识列。 如果你写入表但没有为标识列提供值,它将自动分配一个唯一且统计上递增的值(如果 step 为负数则递减)。 只有 Delta 表支持此子句。 此子句只能用于具有 BIGINT 数据类型的列。

      自动分配的值以 start 开头并以 step 为增量。 分配的值是唯一的,但不保证是连续的。 这两个参数都是可选的,默认值为 1。 step 不能为 0。

      如果自动分配的值超出标识列类型的范围,则查询将失败。

      使用 ALWAYS 时,不能为标识列提供自己的值。

      不支持以下操作:

      • PARTITIONED BY 是标识列
      • UPDATE 是标识列

      注释

      在表上声明标识列会禁用并发事务。 仅在不需要对目标表进行并发写入的用例中使用标识列。

    • 默认default_expression

      适用于:勾选“是” Databricks SQL 勾选“是” Databricks Runtime 11.3 LTS 及更高版本

      为列定义一个 DEFAULT 值,当未指定该列时,将在 INSERT、UPDATE 和 MERGE ... INSERT 上使用该值。

      如果未指定默认值,则 DEFAULT NULL 应用于可为空的列。

      default_expression 可由字面量和内置 SQL 函数或运算符组成,但以下函数除外:

      此外,default_expression 不能包含任何子查询。

      DEFAULT、CSV、JSON 和 PARQUET 源支持 ORC。

    • column_constraint

      将信息主键或信息外键约束添加到流式处理表中的列。

    • MASK 字句

      添加列掩码函数以对敏感数据进行匿名化处理。

      请参阅 行筛选器和列掩码。

    • CONSTRAINT 约束名称 EXPECT (期望表达式) [ ON VIOLATION { 更新失败 | 删除行 } ]

      向流式处理表添加数据质量预期。 这些数据质量预期可以随着时间的推移进行跟踪,并通过流式处理表的事件日志进行访问。 在创建表和刷新表时,FAIL UPDATE 期望会导致处理失败。 如果未满足DROP ROW 预期,则该预期会导致整行被删除。 请参阅通过管道预期管理数据质量。

      expectation_expr 可能包含文本、表中的列标识符以及内置的确定性 SQL 函数或运算符,但以下内容除外:

      此外,expr 不能包含任何子查询。

  • table_constraint

    指定架构时,可以定义主键和外键。 约束具备信息性,系统不会强制执行。 请参阅 SQL 语言参考中的 CONSTRAINT 子句。

    注释

    若要定义表约束,管道必须是启用了 Unity Catalog 的管道。

  • table_clauses

    (可选)指定表的分区、注释和用户定义的属性。 每个子句只能指定一次。

    • 使用 DELTA

      指定数据格式。 唯一的选项是 DELTA。

      此子句是可选的,默认为 DELTA。

    • PARTITIONED BY

      包含一列或多列的可选列表,用于对表进行分区。 与 CLUSTER BY 互斥。

      Liquid 聚类分析提供灵活的优化解决方案进行聚类分析。 请考虑使用 CLUSTER BY 而不是 PARTITIONED BY 用于管道。

    • 按组排序

      对表启用动态聚类,并定义要用作聚类键的列。 使用自动聚类功能与CLUSTER BY AUTO,Databricks智能地选择聚类键以优化查询性能。 与 PARTITIONED BY 互斥。

      请参阅对表使用 liquid 聚类分析。

    • LOCATION

      表数据的可选存储位置。 如果未设置,系统将默认为管道存储位置。

    • 评论

      用于描述表的可选 STRING 文本。

    • TBLPROPERTIES

      表的表属性可选列表。

    • WITH ROW FILTER

    向表中添加行筛选器函数。 将来对该表的查询会收到函数计算结果为 TRUE 的行的子集。 这对于精细的访问控制很有用,因为它允许函数检查调用用户的标识和组成员身份以决定是否筛选某些行。

    请参阅 ROW FILTER 条款。

    • 流

      (可选)定义与表创建内联的 流 。 流是刷新表内容的有状态查询。 如果未 FLOW 指定,则可以改用 AS query 流,也可以单独定义 CREATE FLOW流。 可以指定以类型之一:

      • 按名称插入

        按列名将数据插入表中。 ONCE如果未提供该选项,则查询必须是流式处理查询。 要使用流式处理语义从源中读取,请使用 STREAM 关键字。 如果读取遇到对现有记录的更改或删除,则会引发错误。 从静态源或仅限追加的源读取是最安全的。

        注释

        FLOW INSERT BY NAME 等效于使用 AS query。 以下两个语句具有相同的行为:

        CREATE OR REFRESH STREAMING TABLE raw_data
        AS SELECT * FROM STREAM read_files('abfss://my_path');
        
        CREATE OR REFRESH STREAMING TABLE raw_data
        FLOW INSERT BY NAME SELECT * FROM STREAM read_files('abfss://my_path');
        
      • 一次

        (可选)将流定义为一次性流,例如回填。 提供时 ONCE ,查询不是流式查询,默认情况运行一次。 如果使用完全刷新刷新表,流 ONCE 将再次运行以重新创建数据。 ONCE INSERT BY NAME仅适用于流。

      • AUTO CDC

        Important

        在 Databricks Runtime 17.3 及更高版本和 PREVIEW 管道通道中可用。

        定义处理 AUTO CDC 从源到表中的更改数据捕获(CDC)记录的流。 当源数据包含 CDC 语义时使用 AUTO CDC 。 请参阅 AUTO CDC API:使用管道简化变更数据捕获。

      • REPLACE WHERE 谓词 BY NAME 查询

        定义只 REPLACE WHERE 重新计算和覆盖匹配 predicate的行的流,使所有其他行保持不变。 用于 REPLACE WHERE 对联接和聚合、后期到达的数据、架构演变和回填进行增量批处理。 BY NAME 必需。 请参阅 使用 REPLACE WHERE 流的批处理。

  • AS 查询

    此子句使用 query 中的数据来填充表。 此查询必须是流式处理查询。 要使用流式处理语义从源中读取,请使用 STREAM 关键字。 如果读取遇到对现有记录的更改或删除,则会引发错误。 从静态源或仅限追加的源读取是最安全的。 若要引入具有更改提交的数据,可以添加 skipChangeCommits 读取选项来处理错误。

    同时指定 query 和 table_specification 时,table_specification 中指定的表架构必须包含 query 返回的所有列,否则会出现错误。 在 table_specification 中指定但未由 query 返回的任何列在查询时都返回 null 值。

    有关流数据的详细信息,请参阅 使用管道转换数据。

    • 读取选项

      可以在查询中指定读取选项,以配置从源中读取数据的方式。 例如,可以指定 skipChangeCommits 跳过源数据中的任何更改提交。 读取选项在查询子句中 WITH 指定为映射。 例如:

      SELECT * FROM STREAM source_table WITH (SKIPCHANGECOMMITS=TRUE, STARTINGVERSION=X)
      

      这是 =TRUE 可选的,因此还可以指定如下所示的布尔选项:

      SELECT * FROM STREAM source_table WITH (SKIPCHANGECOMMITS)
      

      注释

      Databricks Runtime 17.3 及更高版本仅支持读取选项。

      Delta 支持以下读取选项,有关每个选项的详细信息,请参阅 Delta Lake 表流式读取和写入。

      • maxFilesPerTrigger
      • maxBytesPerTrigger
      • startingVersion
      • startingTimestamp
      • readChangeFeed
      • withEventTimeOrder
      • skipChangeCommits

所需的权限

管道的运行方式用户必须具有以下权限:

  • 对流式处理表引用的基表的 SELECT 特权。
  • 对父目录的 USE CATALOG 特权和对父架构的 USE SCHEMA 特权。
  • 对流式处理表架构的 CREATE MATERIALIZED VIEW 特权。

为了使用户能够更新在其中定义流式处理表的管道,他们需要:

  • 对父目录的 USE CATALOG 特权和对父架构的 USE SCHEMA 特权。
  • 流式处理表的所有权或对流式处理表的 REFRESH 特权。
  • 流式处理表所有者必须对流式处理表引用的基表具有 SELECT 特权。

要使用户能够查询生成的流式处理表,他们需要:

  • 对父目录的 USE CATALOG 特权和对父架构的 USE SCHEMA 特权。
  • SELECT 流式处理表的特权。

局限性

  • 只有表所有者才能刷新流式处理表以获取最新数据。
  • 流式处理表上不允许 ALTER TABLE 命令。 应该通过 CREATE OR REFRESH 或 ALTER STREAMING TABLE 语句更改该表的定义和属性。
  • 不支持通过 DML 命令(如 INSERT INTO 和 MERGE)来发展表模式。
  • 流式数据表不支持以下命令:
    • CREATE TABLE ... CLONE <streaming_table>
    • COPY INTO
    • ANALYZE TABLE
    • RESTORE
    • TRUNCATE
    • GENERATE MANIFEST
    • [CREATE OR] REPLACE TABLE
  • 不支持重命名表或更改所有者。

例子

-- Define a streaming table from a volume of files:
CREATE OR REFRESH STREAMING TABLE customers_bronze
AS SELECT * FROM STREAM read_files("/databricks-datasets/retail-org/customers/*", format => "csv")

-- Define a streaming table from a streaming source table:
CREATE OR REFRESH STREAMING TABLE customers_silver
AS SELECT * FROM STREAM(customers_bronze)

-- Use automatic liquid clustering to let Databricks choose the clustering columns:
CREATE OR REFRESH STREAMING TABLE customers_bronze_auto
CLUSTER BY AUTO
AS SELECT * FROM STREAM read_files("/databricks-datasets/retail-org/customers/*", format => "csv")

-- Define a table with a row filter and column mask:
CREATE OR REFRESH STREAMING TABLE customers_silver (
  id int COMMENT 'This is the customer ID',
  name string,
  region string,
  ssn string MASK catalog.schema.ssn_mask_fn COMMENT 'SSN masked for privacy'
)
WITH ROW FILTER catalog.schema.us_filter_fn ON (region)
AS SELECT * FROM STREAM(customers_bronze)

-- Define a streaming table with an identity column:
CREATE OR REFRESH STREAMING TABLE customers_with_id (
  customer_id BIGINT GENERATED ALWAYS AS IDENTITY,
  name string,
  region string
)
AS SELECT name, region FROM STREAM(customers_bronze)

-- Define a streaming table that you can add flows into:
CREATE OR REFRESH STREAMING TABLE orders;

-- Define a streaming table with an inline append flow:
CREATE OR REFRESH STREAMING TABLE raw_data
FLOW INSERT BY NAME SELECT * FROM STREAM read_files('abfss://my_path');

-- Define a streaming table with an inline AUTO CDC flow:
CREATE OR REFRESH STREAMING TABLE target
FLOW AUTO CDC
FROM stream(cdc_data.users)
KEYS (userId)
SEQUENCE BY sequenceNum
STORED AS SCD TYPE 1;