Delta Live Tables工作流Schema覆盖及推断疑问咨询
问题背景
我是DLT新手,运行一个简单的数据管道时遇到Schema不匹配错误,相关代码和错误信息如下:
出错的表代码
@dlt.table( table_properties={ "quality" : "silver" } ) def silver_catalog_product(): eav_attributes = dlt.read("eav_attribute"); entityDf = dlt.read("product_entity") entity_datetime_Df = dlt.read("product_entity_datetime") entity_datetime_Df = entity_datetime_Df.join(eav_attributes,entity_datetime_Df.attribute_id == eav_attributes.attribute_id,"inner") entity_datetime_Df = entity_datetime_Df.groupBy("entity_id").pivot("attribute_code").agg(first("value")) df = (entityDf .join(entity_datetime_Df, entityDf.entity_id == entity_datetime_Df.entity_id, "inner") .drop(entity_datetime_Df.entity_id) ) return df
运行时错误信息
To enable schema migration using DataFrameWriter or DataStreamWriter, please set:
'.option("mergeSchema", "true")'.
For other operations, set the session configuration
spark.databricks.delta.schema.autoMerge.enabled to "true". See the documentation
specific to the operation for details.Table schema:
root
-- entity_id: long (nullable = true)
-- entity_type_id: integer (nullable = true)
-- attribute_set_id: integer (nullable = true)
-- type_id: string (nullable = true)
-- sku: string (nullable = true)
-- created_at: timestamp (nullable = true)
-- updated_at: timestamp (nullable = true)
-- has_options: integer (nullable = true)
-- required_options: integer (nullable = true)Data schema:
root
-- entity_id: long (nullable = true)
-- entity_type_id: integer (nullable = true)
-- attribute_set_id: integer (nullable = true)
-- type_id: string (nullable = true)
-- sku: string (nullable = true)
-- created_at: timestamp (nullable = true)
-- updated_at: timestamp (nullable = true)
-- has_options: integer (nullable = true)
-- required_options: integer (nullable = true)
-- custom_design_from: timestamp (nullable = true)
-- custom_design_to: timestamp (nullable = true)
-- news_from_date: timestamp (nullable = true)
-- news_to_date: timestamp (nullable = true)
-- price_update_type_revert_date: timestamp (nullable = true)
-- special_from_date: timestamp (nullable = true)
-- special_to_date: timestamp (nullable = true)To overwrite your schema or change partitioning, please set:
'.option("overwriteSchema", "true")'.Note that the schema can't be overwritten when using
'replaceWhere'.
核心问题
- 熟悉常规Delta PySpark作业里的
overwriteSchema/mergeSchema选项,但不知道怎么在DLT里启用overwriteSchema。 - 为什么作业推断的Schema是
entityDf的Schema,而不是最终返回的df的实际Schema?
解答
问题1:DLT中启用overwriteSchema的方式
DLT不需要用常规的.option配置,有两种实现方式:
- 单表级别启用:直接在
@dlt.table的table_properties中添加delta.writeSchema参数并设为overwrite,修改后的代码如下:
@dlt.table( table_properties={ "quality" : "silver", "delta.writeSchema": "overwrite" } ) def silver_catalog_product(): # 原代码逻辑保持不变
- 全局管道级别启用:如果多个表需要自动处理Schema变更,可在DLT管道的配置页面添加会话参数
spark.databricks.delta.schema.autoMerge.enabled=true,该参数会自动合并新增的Schema字段。
问题2:Schema推断为entityDf的原因
主要有两个核心因素:
- 若该表之前已经创建过(比如试运行时仅返回过
entityDf),DLT会沿用已固化的表Schema,不会自动更新为新的DataFrame Schema。 - 代码中使用了
pivot操作,该操作生成的列是动态的,DLT首次推断Schema时可能未捕获到这些动态列,导致后续返回包含新增列的DataFrame时,与已有表Schema不匹配。
内容的提问来源于stack exchange,提问作者Oliver

