You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.25 17:15:47