使用 Auto Loader 进行自动类型扩展

重要

此功能在 Databricks Runtime 16.4 及更高版本中为公共预览版。

自动加载程序会在新数据文件到达云存储空间时以增量方式高效地对其进行处理。 它还通过自动处理复杂的架构更改来减少管道维护。 例如,可以将自动加载程序配置为自动检测已加载数据的架构,从而允许在不显式声明数据架构的情况下初始化表。 还可以随着新列的引入而改进表架构,无需随时间推移手动跟踪和应用架构更改。 自动加载程序甚至可以在已获救的数据列中拯救出人意料的数据(例如,由于数据类型不同),从而帮助你避免数据丢失。

但是,已获救的数据列要求你手动处理任何数据类型更改。

若要自动处理其中一些数据类型更改,请在自动加载器中使用类型扩大。 Delta Lake 现在支持不同的数据类型扩大更改,而无需数据重写或用户干预。 请参阅 Delta Lake 类型扩展。 模式演进的新模式 addNewColumnsWithTypeWidening,可在数据类型发生兼容性变更时自动演进模式。

可以扩大基元类型,例如intlongfloatdouble等等。 类型扩大适用于自动加载程序中架构演变支持的所有文件格式。 这包括文本格式(如 JSON、CSV 或 XML)和二进制格式(如 Avro 或 Parquet)。 现有架构演变模式(如addNewColumnsrescuefailOnNewColumnsnone)的架构演变行为没有变化。

支持的类型更改

支持以下类型更改:

源类型 支持更多类型
byte shortintlongdecimaldouble
short intlongdecimaldouble
int longdecimaldouble
long decimal
float double
decimal decimal 实现更高精度和更大规模
date timestampNTZ (仅适用于 Parquet 文件)

将任何数值类型加宽为 decimal时,自动加载程序将精度扩大为 decimal 等于或大于起始精度。 如果增加范围,总精度将增加相应程度。

整数类型的起始精度如下:

类型 起始精度
byte 10
short 10
int 10
long 20

例如,如果列的当前类型是 int,并且读取的文件中包含该列类型为 decimal(5, 2),则自动加载程序会将该列的类型扩展为 decimal(12, 2)

先决条件

要在 Auto Loader 中使用类型扩展,必须满足以下要求:

  • 使用 Databricks Runtime 16.4 或更高版本。
  • 如果写入接收端是 Delta Lake 表,请使用以下方法之一为该 Delta Lake 表启用类型扩展:
    • 如果使用现有表:

      ALTER TABLE <table_name> SET TBLPROPERTIES ('delta.enableTypeWidening' = 'true')
      
    • 如果创建的新表已启用类型扩展:

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

有关 Delta Lake 表中类型扩大的详细信息,请参阅 类型扩大

通过架构演变启用类型扩展

若要在使用架构演变时通过 Auto Loader 使用类型扩展,请指定 addNewColumnsWithTypeWidening 。 自动加载程序会在处理数据时检测新列的添加和类型更改。

Python

query = (spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "csv")
  .option("cloudFiles.inferColumnTypes", True)
  .option("cloudFiles.schemaLocation", <schemaPath>)
  .option("cloudFiles.schemaEvolutionMode", "addNewColumnsWithTypeWidening")
  .load(<inputPath>)
  .writeStream
  .option("mergeSchema", "true")
  .option("checkpointLocation", <checkpointPath>)
  .trigger(availableNow=True)
  .toTable("table_name")
)

Scala

val query = spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "csv")
  .option("cloudFiles.inferColumnTypes", true)
  .option("cloudFiles.schemaLocation", <schemaPath>)
  .option("cloudFiles.schemaEvolutionMode", "addNewColumnsWithTypeWidening")
  .load(<inputPath>)
  .writeStream
  .option("mergeSchema", "true")
  .option("checkpointLocation", <checkpointPath>)
  .trigger(Trigger.AvailableNow())
  .toTable("table_name")

当 Auto Loader 检测到类型扩展支持的新列或类型更改时,数据流将因 UnknownFieldException 而停止。 在数据流抛出此错误之前,Auto Loader 会对最新的微批次数据执行模式推断,并通过扩展现有列或将新列合并到模式末尾,将最新模式更新到模式位置。

数据类型更改的架构演变行为

如果要引入包含以下内容的 CSV,自动加载程序会将架构推断为 STRUCT<id INT, name STRING, _rescued_data STRING>

id, name
1, John
2, Mary

目标表如下所示:

id 名字 _rescued_data(已恢复的数据)
1 John Null
2 玛丽 Null

现在,引入另一个 CSV 文件,其中列中的值 idINT 类型更宽:

id, name, age
2147483648, Bob, 25

下表介绍了自动加载程序中具有不同架构演变模式的行为和输出:

模式 支持类型扩展的数据类型变更行为
addNewColumns(默认值) 数据类型不会演变,由于数据类型更改,流不会失败。 包含类型不匹配值的列会被设为 NULL,而这些不匹配的值会被添加到救援数据列中。 数据流在新列处失败。
rescue 架构不会演变,流不会因任何架构更改而失败。 包含类型不匹配值的列会被设为 NULL,而这些不匹配的值会被添加到救援数据列中。
failOnNewColumns 数据类型不会演变,由于数据类型更改,流不会失败。 包含类型不匹配值的列会被设为 NULL,而这些不匹配的值会被添加到救援数据列中。 数据流在新列处失败,且未更新模式。
none 不会使架构演变,将忽略新列,并且除非设置 rescuedDataColumn 选项,否则不会补救数据。 流不会因为架构更改而失败。
addNewColumnsWithTypeWidening 数据流故障。 新列将添加到架构中,支持的数据类型更改将扩大。 不支持的数据类型更改(例如从 int 改为 string)被添加到已获救的数据列中。

示例结果

以下部分显示了引入第二个 CSV 文件后每个架构演变模式的推断架构和值。

addNewColumns

架构id: INT、、name: STRINGage: INT_rescued_data: STRING

id 名字 age _rescued_data(已恢复的数据)
1 John Null Null
2 玛丽 Null Null
Null 鲍勃 25 {"id": 2147483648}

rescue

架构id: INT、、 name: STRING_rescued_data: STRING

id 名字 _rescued_data(已恢复的数据)
1 John Null
2 玛丽 Null
Null 鲍勃 {"age": 25, "id": 2147483648}

failOnNewColumns

架构id: INT、、 name: STRING_rescued_data: STRING

id 名字 _rescued_data(已恢复的数据)
1 John Null
2 玛丽 Null
Null 鲍勃 {"id": 2147483648}

none

模式id: INTname: STRING

id 名字
1 John
2 玛丽
Null 鲍勃

addNewColumnsWithTypeWidening

架构id: BIGINT、、name: STRINGage: INT_rescued_data: STRING

id 名字 age _rescued_data(已恢复的数据)
1 John Null Null
2 玛丽 Null Null
2147483648 鲍勃 25 Null

局限性

  • 使用prefersDecimal时无法将选项false设置为 addNewColumnsWithTypeWidening 。 当指定 addNewColumnsWithTypeWidening 时,prefersDecimal 的默认值是 true
  • datetimestampNTZ 的类型扩展仅支持 Parquet 文件。