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

Synapse Serverless查询更新后Delta Lake分区数据出现重复

问题描述

通过Databricks执行ELT流程,将数据按Year分区存储在Delta Lake中:

  • Databricks内查询数据正常,无重复,总计数为407,421
  • 数据更新后,使用Synapse Serverless通过OPENROWSET指定分区路径创建视图查询时,出现重复数据,总计数翻倍至814,842;但使用未指定分区的Delta Lake外部表查询时结果正常
  • 首次写入数据时无此问题,仅在数据更新后出现

Delta Lake分区数据结构

Delta Lake分区数据

Databricks查询验证

-- 无重复数据返回
select PKCOLUMNS, count(*) from mytable group by PKCOLUMNS having count(*)>1
-- 计数准确:407,421
select count(*) from mytable

Synapse Serverless异常场景

CREATE VIEW MY_TABLE_VIEW AS 
SELECT *, 
results.filepath(1) as [Year]
FROM
OPENROWSET(
BULK 'mytable/Year=*/*.parquet',
DATA_SOURCE = 'DeltaLakeStorage',
FORMAT = 'PARQUET'
) 
WITH(
[param1] nvarchar(4000),
[param2] float,
[PKCOLUMNS] nvarchar(4000)
) AS [results]
GO

-- 查询到重复数据
select PKCOLUMNS, count(*) from MY_TABLE_VIEW
group by PKCOLUMNS
having count(*)>1
GO

-- 计数翻倍:814,842
select count(*) from MY_TABLE_VIEW
原因分析

Delta Lake的ACID事务依赖_delta_log目录下的事务日志文件实现:

  • 数据更新时,Delta Lake不会直接修改原Parquet文件,而是生成新的Parquet文件,同时在事务日志中标记旧文件为"已删除"
  • Databricks读取Delta Lake数据时,会解析_delta_log过滤掉已标记删除的文件,因此无重复
  • Synapse Serverless的OPENROWSET直接扫描指定路径下的所有Parquet文件,不识别Delta Lake的事务日志,会同时读取新旧文件,导致重复数据

而未指定分区的Delta外部表基于Delta Lake格式创建,会自动读取_delta_log过滤无效文件,因此查询结果正常。

解决方案

方案1:创建Delta Lake格式的外部表(推荐)

直接基于Delta Lake根目录创建外部表,而非扫描Parquet文件:

-- 先创建Delta格式的文件格式(如果未创建)
CREATE EXTERNAL FILE FORMAT DeltaFormat
WITH (
    FORMAT_TYPE = DELTA
);

-- 创建Delta外部表
CREATE EXTERNAL TABLE my_delta_external_table
(
    param1 nvarchar(4000),
    param2 float,
    PKCOLUMNS nvarchar(4000),
    Year int -- 分区列
)
WITH (
    LOCATION = 'mytable/', -- 指向Delta Lake根目录
    DATA_SOURCE = 'DeltaLakeStorage',
    FILE_FORMAT = DeltaFormat
);

-- 基于外部表创建视图
CREATE VIEW MY_TABLE_VIEW AS
SELECT * FROM my_delta_external_table;

此方式自动识别Delta事务日志,过滤已删除文件,确保查询结果准确。

方案2:执行Delta Lake的VACUUM清理旧文件

在Databricks中执行VACUUM操作,物理删除已标记为删除的Parquet文件:

-- 清理所有超过0小时的旧文件(根据业务需求调整保留时间)
VACUUM mytable RETAIN 0 HOURS;

注意:VACUUM会清除旧版本数据,影响Delta Lake的**时间旅行(Time Travel)**功能,需根据业务场景谨慎使用。

方案3:通过解析_delta_log过滤有效文件(不推荐)

手动解析_delta_log中的事务记录,筛选出当前有效的Parquet文件后再通过OPENROWSET读取。此方式实现复杂、维护成本高,仅在无法使用前两种方案时考虑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 22:13:13