類型擴展到更大範圍

在 Databricks Runtime 15.4 LTS 及以上版本的 Delta Lake 資料表中,類型寬度允許你將欄位資料型態變更為更寬的型別,而無需重寫資料檔案。

所有 Unity 目錄受控數據表預設都會使用 Delta Lake。 請參閱 Unity Catalog 管理的 Delta Lake 與 Apache Iceberg 表格。

Note

啟用型別擴展會升級讀取器和寫入器通訊協定。 這可能會影響與外部 Delta Lake 用戶端的相容性。 請參閱 Delta Lake 功能相容性和通訊協定。

已啟用型別擴展的資料表僅能由 Databricks Runtime 15.4 LTS 及更新版本讀取。

支援的型別變更

您可以根據下列規則來擴大類型:

來源類型 支援更廣泛的類型
BYTE SHORT、、 INT、 BIGINT、 DECIMAL、 DOUBLE
SHORT INT、BIGINT、DECIMAL、DOUBLE
INT BIGINT、DECIMAL、DOUBLE
BIGINT DECIMAL
FLOAT DOUBLE
DECIMAL 具有更高精確度和規模的 DECIMAL
DATE TIMESTAMP_NTZ
VOID 任何類型

最上層數據行和巢狀結構、地圖和陣列內的欄位支援類型變更。

Note

VOID 對任何類型都不需要在表格上啟用型別拓寬。 任何執行更新 VOID 欄位類型的操作均可成功,無需額外設定。 VOID 型別擴展功能可在 Databricks Runtime 18.2 及以上版本使用。

小數處理方式

當操作將整數類型提升為 decimal 或 double 且下游將該值寫回整數欄位時,Spark 預設會截斷數值的小數部分。 關於指派政策行為的詳細資訊,請參見 「儲存分配」。

將任何數值類型變更為 decimal時,總有效位數必須等於或大於起始有效位數。 如果您也增加規模,則總精確度必須增加相應的數量。

byte、short和 int 類型的最小目標為 decimal(10,0)。 long 最低目標為 decimal(20,0)。

如果您要將兩個小數位數新增至具有 decimal(10,1)的欄位,則最小目標為 decimal(12,3)。

啟用類型拓寬

Note

啟用型別擴展會升級讀取器和寫入器通訊協定。 這可能會影響與外部 Delta Lake 用戶端的相容性。 請參閱 Delta Lake 功能相容性和通訊協定。

您可以將 table 屬性設定 delta.enableTypeWidening 為 true,以在現有的資料表上啟用類型擴大:

  ALTER TABLE <table_name> SET TBLPROPERTIES ('delta.enableTypeWidening' = 'true')

您也可以在資料表建立期間啟用類型擴大:

  CREATE TABLE T(c1 INT) TBLPROPERTIES('delta.enableTypeWidening' = 'true')

手動套用類型變更

使用ALTER COLUMN命令手動變更類型:

ALTER TABLE <table_name> ALTER COLUMN <col_name> TYPE <new_type>

這項作業會更新數據表架構,而不需重寫基礎數據檔。 如需詳細資訊,請參閱 ALTER TABLE。

使用自動架構演進擴大類型

利用 schema 演化並擴大型別,更新目標資料表中的資料型態,使其與輸入資料類型相符。

Note

若未啟用類型擴大,架構演進一律會嘗試向下轉換數據,以符合目標數據表中的數據行類型。 如果你不想自動擴大目標資料表的資料型態,必須在啟用 schema 演化的工作負載前關閉型別擴大功能。

若要在資料匯入期間使用架構演進來調整欄位的數據類型,您必須符合下列條件:

  • 寫入命令會在啟用自動架構演進時執行。
  • 目標資料表已啟用類型廣化。
  • 源數據行類型比目標數據行類型寬。
  • 類型擴展允許類型變更。

不符合所有這些條件的類型不匹配會遵循通常的架構執行規則。 請參閱強制執行結構模式。

範例

以下範例展示了型別拓寬如何與結構演化相關工作。

Python

建立一個具有 INT 欄位的目標資料表,以及一個具有 BIGINT 欄位的來源資料表:

spark.sql("CREATE TABLE target_table (id INT, data STRING) TBLPROPERTIES ('delta.enableTypeWidening' = 'true')")
spark.sql("CREATE TABLE source_table (id BIGINT, data STRING)")

在附加過程中,利用 saveAsTable() 結構演化自動將 INT 欄位擴大為 BIGINT :

spark.table("source_table").write.mode("append").option("mergeSchema", "true").saveAsTable("target_table")

與圖式演化的應用 MERGE INTO :

from delta.tables import DeltaTable

source_df = spark.table("source_table")
target_table = DeltaTable.forName(spark, "target_table")

(target_table.alias("target")
  .merge(source_df.alias("source"), "target.id = source.id")
  .withSchemaEvolution()
  .whenMatchedUpdateAll()
  .whenNotMatchedInsertAll()
  .execute()
)

Scala

建立一個具有 INT 欄位的目標資料表,以及一個具有 BIGINT 欄位的來源資料表:

spark.sql("CREATE TABLE target_table (id INT, data STRING) TBLPROPERTIES ('delta.enableTypeWidening' = 'true')")
spark.sql("CREATE TABLE source_table (id BIGINT, data STRING)")

在附加過程中,利用 saveAsTable() 結構演化自動將 INT 欄位擴大為 BIGINT :

spark.table("source_table").write.mode("append").option("mergeSchema", "true").saveAsTable("target_table")

與圖式演化的應用 MERGE INTO :

import io.delta.tables.DeltaTable

val sourceDf = spark.table("source_table")
val targetTable = DeltaTable.forName(spark, "target_table")

targetTable.alias("target")
  .merge(sourceDf.alias("source"), "target.id = source.id")
  .withSchemaEvolution()
  .whenMatched().updateAll()
  .whenNotMatched().insertAll()
  .execute()

SQL

建立一個具有 INT 欄位的目標資料表,以及一個具有 BIGINT 欄位的來源資料表:

CREATE TABLE target_table (id INT, data STRING) TBLPROPERTIES ('delta.enableTypeWidening' = 'true');
CREATE TABLE source_table (id BIGINT, data STRING);

在附加過程中,利用 INSERT INTO 結構演化自動將 INT 欄位擴大為 BIGINT :

INSERT WITH SCHEMA EVOLUTION INTO target_table SELECT * FROM source_table;

與圖式演化的應用 MERGE INTO :

MERGE WITH SCHEMA EVOLUTION INTO target_table
USING source_table
ON target_table.id = source_table.id
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *;

自動加載器

Important

自動載入器中的型別擴展支援目前已在公開預覽階段。

自動載入器支援透過自動結構演化來擴展型別。 當你使用 Auto Loader 將資料匯入 Delta Lake 資料表,並啟用類型寬度和結構演化時,欄位類型會自動被擴大以匹配輸入資料。

(spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("cloudFiles.schemaLocation", "<path-to-schema-location>")
  .load("<path-to-source-data>")
  .writeStream
  .option("mergeSchema", "true")
  .option("checkpointLocation", "<path-to-checkpoint>")
  .trigger(availableNow=True)
  .toTable("table_name")
)

請參見 自動型別擴展與 Auto Loader。 此外,目標表必須啟用型別拓寬。 請參見啟用類型擴展。

停用類型擴展表格功能

您可以將 屬性設定為 false,以防止在啟用的資料表上意外擴大類型:

  ALTER TABLE <table_name> SET TBLPROPERTIES ('delta.enableTypeWidening' = 'false')

此設定可防止未來對表格的型別變更,但不會移除型別擴大表的功能或還原先前的型別變更。

如果您需要完全移除類型擴展的表格功能,您可以使用 DROP FEATURE 命令,如下列範例所示:

 ALTER TABLE <table-name> DROP FEATURE 'typeWidening' [TRUNCATE HISTORY]

Note

對於使用 Databricks Runtime 15.4 LTS 啟用型別擴展的資料表,您必須改為移除此功能 typeWidening-preview。

在移除型別擴展時,Databricks 會重寫所有不符合目前資料表模式的資料檔案。 請參閱 刪除 Delta Lake 資料表功能和降級資料表協定。

從 Delta Lake 資料表串流讀取資料

結構化串流中對類型擴展的支援可在 Databricks Runtime 16.4 LTS 及以上版本中提供。

當從已啟用型別擴寬的 Delta Lake 資料表進行串流時,你可以在目標資料表上使用 mergeSchema 選項啟用結構描述演進,為串流查詢設定自動型別擴寬。 目標表必須啟用資料型別擴展功能。 請參見啟用類型擴展。

Python

(spark.readStream
  .table("delta_source_table")
  .writeStream
  .option("checkpointLocation", "/path/to/checkpointLocation")
  .option("mergeSchema", "true")
  .toTable("output_table")
)

Scala

spark.readStream
  .table("delta_source_table")
  .writeStream
  .option("checkpointLocation", "/path/to/checkpointLocation")
  .option("mergeSchema", "true")
  .toTable("output_table")

當 mergeSchema 啟用且目標資料表啟用類型擴展時:

  • 型別變更會自動套用到下游資料表,無需人工介入。
  • 新的欄位會自動加入到下游資料表架構中。

未 mergeSchema 啟用時,值會依 spark.sql.storeAssignmentPolicy 設定處理,預設會下播與目標欄位類型相符的值。 欲了解更多關於指派政策行為的資訊,請參見 「儲存分配」。

處理串流中的類型變更

從 Delta Lake 資料表串流時,你可以提供結構追蹤位置,追蹤非加法結構變更,包括型別變更。 在 Databricks Runtime 18.0 及以下版本中,提供結構追蹤位置是必要的,而在 Databricks Runtime 18.1 及以上版本則為可選。

你不能用 SQL 設定 a schemaTrackingLocation 。 請參見 不支援的功能。

schemaTrackingLocation 必須設定在與串流檢查點相同路徑內的位置。 例如:

Python

checkpoint_path = "/path/to/checkpointLocation"

(spark.readStream
  .option("schemaTrackingLocation", checkpoint_path)
  .table("delta_source_table")
  .writeStream
  .option("checkpointLocation", checkpoint_path)
  .toTable("output_table")
)

Scala

val checkpointPath = "/path/to/checkpointLocation"

spark.readStream
  .option("schemaTrackingLocation", checkpointPath)
  .table("delta_source_table")
  .writeStream
  .option("checkpointLocation", checkpointPath)
  .toTable("output_table")

設定模式追蹤位置後,串流在偵測到型別改變時演化其追蹤結構,然後停止。 此時,你必須處理型別變更,例如在下游資料表上啟用型別擴大,或更新串流查詢。

要繼續處理,請設定 Spark 設定 spark.databricks.delta.streaming.allowSourceColumnTypeChange 或 DataFrame 讀取器選項 allowSourceColumnTypeChange,如下範例所示:

Python

checkpoint_path = "/path/to/checkpointLocation"

(spark.readStream
  .option("schemaTrackingLocation", checkpoint_path)
  .option("allowSourceColumnTypeChange", "<delta_source_table_version>")
  # alternatively to allow all future type changes for this stream:
  # .option("allowSourceColumnTypeChange", "always")
  .table("delta_source_table")
  .writeStream
  .option("checkpointLocation", checkpoint_path)
  .toTable("output_table")
)

Scala

val checkpointPath = "/path/to/checkpointLocation"

spark.readStream
  .option("schemaTrackingLocation", checkpointPath)
  .option("allowSourceColumnTypeChange", "<delta_source_table_version>")
  // alternatively to allow all future type changes for this stream:
  // .option("allowSourceColumnTypeChange", "always")
  .table("delta_source_table")
  .writeStream
  .option("checkpointLocation", checkpointPath)
  .toTable("output_table")

SQL

  -- To unblock for this particular stream just for this series of schema change(s):
  SET spark.databricks.delta.streaming.allowSourceColumnTypeChange.ckpt_<checkpoint_id> = "<delta_source_table_version>"
  -- To unblock for this particular stream:
  SET spark.databricks.delta.streaming.allowSourceColumnTypeChange = "<delta_source_table_version>"
  -- To unblock for all streams:
  SET spark.databricks.delta.streaming.allowSourceColumnTypeChange = "always"

當串流停止時,會顯示錯誤訊息,其中包含檢查點 ID <checkpoint_id> 和 Delta Lake 來源資料表版本 <delta_source_table_version>。

欲了解完整的串流三角洲湖選項列表,請參見 三角洲湖。

湖流量管線

你可以在管線層級或針對個別資料表,為 Lakeflow 管線啟用型別擴展。 型別寬度允許在管線執行時自動擴大欄位類型,而無需完整刷新串流資料表。 實體化檢視的類型變更總是會觸發完整重新計算,當類型變更應用於來源資料表時,依賴該資料表的實體化檢視則需全面重新計算以反映新類型。

啟用整條管線的類型擴展

要啟用管線中所有資料表的型別拓寬,請設定管線配置 pipelines.enableTypeWidening:

JSON

{
  "configuration": {
    "pipelines.enableTypeWidening": "true"
  }
}

YAML

configuration:
  pipelines.enableTypeWidening: 'true'

啟用特定資料表的型別擴展

你也可以透過設定 table 屬性 delta.enableTypeWidening來啟用個別資料表的型別擴展:

Python

import dlt

@dlt.table(
  table_properties={"delta.enableTypeWidening": "true"}
)
def my_table():
  return spark.readStream.table("source_table")

SQL

CREATE OR REFRESH STREAMING TABLE my_table
TBLPROPERTIES ('delta.enableTypeWidening' = 'true')
AS SELECT * FROM source_table

與下游讀取器的相容性

已啟用類型擴大的數據表只能在 Databricks Runtime 15.4 LTS 和更新版本讀取。 如果你希望在 Databricks Runtime 14.3 及以下版本中,讀者能讀取管線中已啟用型別擴大的表格,則必須執行以下其中一項:

  • 移除 delta.enableTypeWidening/pipelines.enableTypeWidening 屬性或將其設為 false,以關閉型別擴展,並觸發資料表的完整重新整理。
  • 在你的桌面上啟用 相容模式 。

開放共享

Note

Databricks 執行環境 16.1 及更新版本開始提供 OpenSharing 的型別擴大支援。

Databricks-to-Databricks OpenSharing 支援共用已啟用型別擴展的 Delta Lake 資料表。 提供者和收件者必須位於 Databricks Runtime 16.1 或更新版本。

若要使用 OpenSharing 從已啟用型別擴展的 Delta Lake 資料表讀取變更資料饋送,您必須將回應格式設為 delta:

spark.read
  .format("deltaSharing")
  .option("responseFormat", "delta")
  .option("readChangeFeed", "true")
  .option("startingVersion", "<start version>")
  .option("endingVersion", "<end version>")
  .load("<table>")

不支援跨類型變更讀取變更資料流。 您必須改為將作業分割成兩個不同的讀取,一個結束於包含類型變更的數據表版本,另一個從包含類型變更的版本開始。

Limitations

Apache Iceberg 兼容性

Apache Iceberg 並不支援型別擴大所涵蓋的所有型別變更。 參見 冰山圖式演化。

不受支援的類型變更包括以下幾項:

  • byte、short、int、long 至 decimal 或 double
  • 小數範圍增加
  • date 到 timestampNTZ

當你在 Delta Lake 表格啟用 Iceberg 讀取時,套用上述類型變更會產生錯誤。 請參閱 使用 Iceberg 用戶端讀取 Delta Lake 資料表。

如果你對 Delta Lake 表格套用這些不支援的類型變更,你有兩個選項:

  • 重新產生 Iceberg 元資料:使用以下指令在不使用類型擴大表格功能的情況下重新生成 Iceberg 元資料:

    ALTER TABLE <table-name> SET TBLPROPERTIES ('delta.universalFormat.config.icebergCompatVersion' = '<version>')
    

    這讓你在套用不相容的型別變更後,仍能維持與 Iceberg 讀取的相容性。

  • 移除型別寬寬表功能:詳見 「停用型別寬寬表功能」。

型別相關函數

有些 SQL 函式會回傳依輸入資料型別而異的結果。 例如, hash 函數 對於同一邏輯值,若參數類型不同,會回傳不同的雜湊值: hash(1::INT) 回傳的結果 hash(1::BIGINT)與 不同。

其他型別相關函數有:xxhash64、bit_get、bit_reversetypeof、。

若要在使用這些函式的查詢中獲得穩定結果,必須明確將值鑄造成所需的類型:

Python

spark.read.table("table_name") \
  .selectExpr("hash(CAST(column_name AS BIGINT))")

Scala

spark.read.table("main.johan_lasperas.dlt_type_widening_bronze2")
  .selectExpr("hash(CAST(a AS BIGINT))")

SQL

-- Use explicit casting for stable hash values
SELECT hash(CAST(column_name AS BIGINT)) FROM table_name

不支援的功能

  • 當從具有型別變更的 Delta Lake 資料表進行串流時,您無法使用 SQL 設定結構描述追蹤位置。
  • 使用 OpenSharing 時,您無法將已啟用型別擴展的資料表分享給非 Databricks 取用者。