如何让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
相关产品推荐
相关产品推荐

