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

如何让Delta Live Table仅增量加载ADLS中的新数据?

问题描述

现有运行在Delta Live Table(DLT)管道中的非流式代码,每4小时调度运行,但每次全量加载ADLS存储容器中的所有数据,随着数据量增长,运行时长从30分钟增至3小时30分钟。希望修改为仅处理新增记录/文件并追加到现有表,尝试用readStream+cloudFiles时出现错误。

现有代码

def generate_table(live_table, landing_file_path):
    file_path = landing_file_path
    table_name = live_table.replace("-", "_").replace(".", "_").split('=')[1]
    cn=table_name.split(f'{table_name_prefix.lower()}')[1]
    category_name="_".join(cn.split("_", 2)[:2])
    table_folder_path="_".join(cn.split("_")[2:])

    @dlt.table(
        name= f"{kafka_cmd_table_prefix}{table_name}",
        path= f"{landing_file_path}/{category_name}/{table_folder_path}",
        comment="Landing zone (Raw) data capture for " + str(table_name)
    )
    def create_live_table():
        pqt_path = f"{pre_landing_file_path}/{live_table}"
        return (
          (spark.read.format("delta").option("ignoreMissingFiles", "true").load(pqt_path))
        )


client_list = dbutils.fs.ls(pre_landing_file_path)
for c in client_list:
    if c.name[:-1].startswith(f"{topic_prefix}"):
        for x in dbutils.fs.ls(c[0]):
            if folder_name == c.name[:-1]:
                print("same path found !!!")
            else:
                # if not folder_name.startswith(("landing", "raw", "processed")):
                if x[0].endswith('.parquet'):
                    folder_name = c.name[:-1]
                    if not folder_name.startswith(("landing", "raw", "procesed")):
                        generate_table(c.name[:-1], landing_file_path)

尝试的修改及错误

修改后的create_live_table函数:

def create_live_table():
        pqt_path = f"{pre_landing_file_path}/{live_table}"
        return (
        #   (spark.read.format("delta").option("ignoreMissingFiles", "true").load(pqt_path))
            (spark.readStream.format("cloudFiles").option("cloudFiles.format", "parquet").load(pqt_path))
        )

错误信息:

A transaction log for Delta was found at `abfss://.../_delta_log`,
but you are trying to read from `abfss://..` using format("cloudFiles"). You must use
'format("delta")' when reading and writing to a delta table.

疑问:是否可以中途从delta格式切换为cloudFiles?寻求解决方法。


解决方案

1. 不能直接切换为cloudFiles读取Delta目录

当路径下存在Delta事务日志(_delta_log)时,必须用format("delta")读取,不能用cloudFiles。cloudFiles仅用于读取未被Delta管理的原始文件(如独立Parquet、CSV文件),而非Delta表目录。

2. 实现增量加载的可行方案

方案一:利用Delta表版本实现增量读取

如果源路径是Delta表,可通过版本追踪新增数据:

@dlt.table(...)
def create_live_table():
    target_table_full_name = f"{kafka_cmd_table_prefix}{table_name}"
    try:
        # 获取目标表的最新版本
        target_version = spark.sql(f"DESCRIBE HISTORY {target_table_full_name}").select("version").first()[0]
        # 读取源表中比目标表版本新的数据
        return spark.read.format("delta").option("ignoreMissingFiles", "true")\
            .option("startingVersion", target_version + 1)\
            .load(pqt_path)
    except Exception:
        # 目标表不存在时全量加载
        return spark.read.format("delta").option("ignoreMissingFiles", "true").load(pqt_path)

注意:该方案要求源Delta表的版本连续递增,且能通过版本准确追踪新增数据。

方案二:基于时间戳的增量筛选

如果表中有记录时间戳字段(如ingestion_time),可通过时间范围过滤增量数据:

@dlt.table(...)
def create_live_table():
    target_table_full_name = f"{kafka_cmd_table_prefix}{table_name}"
    try:
        # 获取目标表的最新记录时间
        max_time = spark.sql(f"SELECT MAX(ingestion_time) FROM {target_table_full_name}").first()[0]
        # 读取源表中时间戳晚于max_time的数据
        return spark.read.format("delta").option("ignoreMissingFiles", "true")\
            .load(pqt_path)\
            .filter(f"ingestion_time > '{max_time}'")
    except Exception:
        # 目标表不存在时全量加载
        return spark.read.format("delta").option("ignoreMissingFiles", "true").load(pqt_path)

若源数据无时间戳字段,可通过文件修改时间筛选新增文件:

def get_new_files(base_path, last_run_ts):
    new_files = []
    for item in dbutils.fs.ls(base_path):
        if item.path.endswith(".parquet") and item.modificationTime > last_run_ts:
            new_files.append(item.path)
    return new_files

@dlt.table(...)
def create_live_table():
    # 从元数据表获取上次运行的时间戳(需提前创建元数据表存储状态)
    last_run_ts = spark.sql("SELECT last_run_time FROM metadata_table WHERE table_name = 'xxx'").first()[0]
    new_files = get_new_files(pqt_path, last_run_ts)
    
    if new_files:
        return spark.read.format("parquet").load(new_files)
    else:
        # 无新增数据时返回空DataFrame,避免全量加载
        empty_schema = spark.read.format("delta").load(pqt_path).schema
        return spark.createDataFrame([], schema=empty_schema)

方案三:切换为流式DLT管道(业务允许时)

如果可以将管道改为流式,直接读取Delta表的变更流,自动实现增量:

@dlt.table(...)
def create_live_table():
    return (
        spark.readStream.format("delta")
            .load(pqt_path)
    )

# 明确指定追加模式(可选,DLT默认会处理)
dlt.append_flow(create_live_table(), target=f"{kafka_cmd_table_prefix}{table_name}")

流式管道会自动追踪Delta表的新增数据,无需手动维护增量状态。

3. 关键注意点

  • 源路径是Delta表时,必须用format("delta")读取,cloudFiles不适用。
  • 非流式管道的增量逻辑需手动维护状态(版本、时间戳、文件修改时间),建议用元数据表存储这些状态信息。
  • 流式DLT管道是增量处理Delta表最省心的方案,适合持续写入的数据源。

内容的提问来源于stack exchange,提问作者Yuva

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 07:20:03