在Databricks DLT任务中合并两个Gold Delta Live Table的问题
解决Delta Live Table中PySpark合并两个Gold表的问题
问题背景
运行Delta Live Table(DLT)任务处理两个独立CSV文件,各自搭建了Bronze→Silver→Gold的数据流链路,得到两个独立的Gold表。原本想用SQL通过%sql魔法命令合并这两个表,但在DLT Python Notebook中运行时触发报错:
Magic commands (e.g. %py, %sql and %run) are not supported with the exception of %pip within a Python notebook. Cells containing magic commands are ignored. Unsupported magic commands were found in the following notebooks
同时误以为PySpark无法实现建表操作,现需要用PySpark完成两个Gold表基于共享外键的合并,生成包含指定列的最终表。
解决方案
1. 核心合并逻辑(DLT中创建合并后的Gold表)
利用DLT的@dlt.table装饰器,直接读取两个已有的Gold表,通过共享外键执行Join操作,筛选所需列后生成新的合并Gold表。以下是示例代码:
假设两个Gold表分别为Gold_or和Gold_another,共享外键为CustomerID,需要合并后保留双方的指定列:
from pyspark.sql import functions as F from pyspark.sql.functions import desc @dlt.table( name="Gold_merged", comment="Merged gold table from two retail datasets joined on CustomerID", table_properties={"quality": "gold"} ) def Gold_merged(): # 读取两个Gold表,若需增量处理则改用dlt.read_stream gold_dataset1 = dlt.read("Gold_or") gold_dataset2 = dlt.read("Gold_another") # 根据业务需求选择Join类型(inner/left/right/full),这里以内连接为例 merged_data = gold_dataset1.join( gold_dataset2, on="CustomerID", # 共享外键 how="inner" ).select( # 按需选择需要保留的列,注意列名重复时需指定表别名 gold_dataset1.UnitPrice.alias("Or_UnitPrice"), gold_dataset1.Quantity.alias("Or_Quantity"), gold_dataset2.ProductCategory.alias("Another_Category"), gold_dataset2.PurchaseDate.alias("Another_PurchaseDate"), "CustomerID" ).sort(desc("CustomerID")) return merged_data
2. 关键注意事项
- 增量处理适配:如果两个源Gold表是增量更新的,将
dlt.read替换为dlt.read_stream,确保合并操作也能实时处理新数据。 - Join类型选择:根据业务场景选择合适的Join方式:
inner:仅保留两边都存在外键的数据left:保留第一个表的所有数据,匹配第二个表的对应数据right/full:分别保留第二个表或两边所有数据
- 列名冲突处理:若两个表存在同名非外键列,需用
alias重命名,避免查询歧义。
3. 修正Silver表的缺失逻辑
你当前的Silver表仅执行了表创建操作,未写入清洗后的数据,需补充数据写入逻辑:
方式一:直接返回清洗视图数据
@dlt.table( name="Silver_or", comment="Clean, merged retail sales data", table_properties={ "quality": "silver", "pipelines.cdc.tombstoneGCThresholdInSeconds": "2" } ) def Silver_or(): return dlt.read_stream("Bronze_or_clean_v")
方式二:CDC合并(适用于需要增量更新去重的场景)
dlt.create_target_table( name="Silver_or", comment="Clean, merged retail sales data", table_properties={ "quality": "silver", "pipelines.cdc.tombstoneGCThresholdInSeconds": "2" } ) dlt.apply_changes( target="Silver_or", source="Bronze_or_clean_v", keys=["id"], # 表主键 sequence_by="LoadDate", # 用于CDC排序的时间字段 except_column_list=["inputFileName"] # 排除不需要同步的列 )
4. 澄清误解:PySpark完全支持建表
DLT中PySpark建表的两种常用方式:
- 用
@dlt.table装饰器定义表,函数返回DataFrame即可自动创建并写入数据 - 用
dlt.create_target_table先创建空表,再通过dlt.apply_changes或写入逻辑填充数据
内容的提问来源于stack exchange,提问作者Anton Kopti
相关产品推荐
相关产品推荐

