CREATE STREAMING TABLE (管道)

串流數據表是支援串流或累加數據處理的數據表。 串流資料表由管線支援。 每次重新整理串流數據表時,新增至源數據表的數據都會附加至串流數據表。 您可以手動或排程更新串流數據表。

若要深入瞭解如何執行或排程重新整理,請參閱 執行管線更新。

語法

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 } |
  { REPLACE USING ( column_name [, ...] ) SEQUENCE BY sequence_column BY NAME query } }

參數

  • REFRESH

    若指定,則建立資料表,或更新現有資料表及其內容。

  • 私人

    建立私人串流數據表。

    • 它們不會新增至目錄中,而且只能在定義的管道內訪問。
    • 它們的名稱可以與目錄中的現有物件相同。 在管線中,如果私有串流表與目錄中的物件名稱相同,該名稱的引用會解析到私有串流表。
    • 私人串流數據表只會在管線的存留期間持續存在,而非僅限於單次更新。

    先前已使用 TEMPORARY 參數建立私人串流數據表。

  • table_name

    新建立數據表的名稱。 完整合格的資料表名稱必須具有唯一性。

  • 表格規格

    這個選擇性子句會定義欄位清單、其類型、屬性、描述和欄位條件約束。

    • column_identifier

      數據行名稱必須是唯一的,且對應至查詢的輸出數據行。

    • 欄位類型

      指定數據行的數據類型。 串流數據表不支援 Azure Databricks 支援的所有資料類型。

    • column_comment

      描述欄的可選 STRING 常值。 此選項必須與 column_type 一起指定。 如果未指定數據行類型,則會略過數據行批注。

    • 產生的 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 預設表達式

      適用範圍:勾選為「是」 Databricks SQL 勾選為「是」 Databricks Runtime 11.3 LTS 和其以上版本

      當數據行未指定時,定義 DEFAULT 值以用於 INSERT、UPDATE和 MERGE ... INSERT 上的數據行。

      如果未指定預設值,DEFAULT NULL 會套用於可為 Null 的欄。

      default_expression 可能由常值和內建 SQL 函式或運算子組成,但下列函式除外:

      default_expression 也不得包含任何子查詢。

      DEFAULT 支援 CSV、JSON、PARQUET 和 ORC 來源。

    • column_constraint

      在串流資料表的欄位中新增資訊主鍵或資訊外鍵限制。

    • MASK 子句

      新增數據行遮罩函式來匿名敏感數據。

      請參閱 行篩選和列遮罩。

    • CONSTRAINT expectation_name EXPECT (expectation_expr) [ON VIOLATION {FAIL UPDATE | DROP ROW}]

      在串流表中新增資料品質的期望值。 這些數據質量預期可以隨著時間追蹤,並透過串流數據表 的事件記錄檔來存取。 FAIL UPDATE 的預期行為會導致在建立或重新整理數據表時,處理過程失敗。 如果不符合預期,整個資料列會因這個DROP ROW預期而被移除。 請參閱 使用管線期望來管理資料品質。

      expectation_expr 可能包含常值、數據表中的數據行標識符,以及決定性的內建 SQL 函式或運算符,但除外:

      expr 也不得包含任何子查詢。

  • 資料表限制

    指定架構時,您可以定義主鍵和外鍵。 條件約束是參考性的,不會強制執行。 請參閱 SQL 語言參考中的 CONSTRAINT 子句。

    備註

    若要定義資料表約束條件,您的管線必須是已啟用 Unity 目錄的管線。

  • 表格條款

    選擇性地指定數據表的數據分割、批注和用戶定義屬性。 每個次子句只能指定一次。

    • 使用 DELTA

      指定數據格式。 唯一的選項是 DELTA。

      這個子句是選填的,預設值為 DELTA。

    • 分區依據

      用於資料表分割的可選一個或多個欄位清單。 與 CLUSTER BY互斥。

      Liquid clustering 提供彈性且優化的叢集解決方案。 請考慮使用CLUSTER BY替代PARTITIONED BY以用於管線。

    • CLUSTER BY

      在數據表上啟用液體叢集,並定義要當做叢集索引鍵使用的數據行。 使用 CLUSTER BY AUTO 的自動叢集功能,Databricks 會智能地選擇叢集索引鍵,以優化查詢效能。 與 PARTITIONED BY互斥。

      請參閱 針對數據表使用液體叢集。

    • 位置

      數據表數據的選擇性儲存位置。 若未設定,系統會預設為管線儲存位置。

    • 評論

      描述表格的可選 STRING 常值。

    • TBLPROPERTIES

      可選的資料表屬性清單。

    • 與 ROW FILTER

    將數據列篩選函式加入至數據表。 該數據表的未來查詢會接收其中函式評估為 TRUE 的資料列子集。 這適用於精細訪問控制,因為它可讓函式檢查叫用使用者的身分識別和群組成員資格,以決定是否篩選特定數據列。

    請參閱 ROW FILTER 條款。

    • 流程

      可選擇性地定義與資料表建立相關的 流程 。 流程是一種有狀態的查詢,用來刷新資料表的內容。 若FLOW未指定,則可改用AS query,或分別定義流。CREATE FLOW 您可以指定以下其中一種流程類型:

      • INSERT 依名稱

        依欄位名稱將資料插入資料表。 若未提供該 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 再次執行以重建資料。 ONCE 只適用於 INSERT BY NAME 流量。

      • AUTO CDC

        這很重要

        可在 Databricks 執行環境 17.3 及以上版本及 PREVIEW Pipelines 通道中取得。

        定義一個 AUTO CDC 流程,將變更資料擷取(CDC)記錄從來源處理到資料表。 當來源資料包含 CDC 語意時使用 AUTO CDC 。 請參閱 AUTOTO CDC API:使用管線簡化變更資料擷取。

      • 以名稱替換 WHERE謂詞 查詢

        定義一個 REPLACE WHERE 流程,只重新計算並覆蓋與 相符 predicate的列,其他列則不受影響。 用於 REPLACE WHERE 增量批次處理連接與聚合、晚到資料、結構演化及回填。 BY NAME 是必要的。 請參見 使用 REPLACE WHERE 流程的批次處理。

      • 取代使用( column_name [, ...])依sequence_column排序 依名稱查詢

        這很重要

        這項功能位於 測試版 (Beta) 中。 需要 Databricks Runtime 18.2 及以上版本。

        定義一個 REPLACE USING 流程,替換所有符合指定鍵欄位的列,並保留其他列不變。 當你的來源是一系列以欄位鍵入的部分快照時,請使用 REPLACE USING 。 SEQUENCE BY 指令更新順序,使得金鑰的最高序列勝出,即使更新順序不正確。 來源必須是串流來源。 BY NAME 是必要的。 請參見部分快照替換,並用替換使用流程。

  • 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 語句來改變。
  • 不支援透過像 INSERT INTO和 MERGE 這樣的 DML 命令來修改資料表結構。
  • 串流資料表不支援下列命令:
    • 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;

-- Define a streaming table with an inline REPLACE USING flow that keeps the latest
-- row for each payment_id:
CREATE OR REFRESH STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);