通过Databricks Auto Loader读取Delta Share中Delta表的流式处理问题
问题描述
通过Delta Share读取其他环境的Delta表,初始流管道运行正常,但当GCP端源表执行重写类更新操作(如CREATE OR REPLACE TABLE)后,现有代码报错,必须重置检查点才能继续,但重置会导致重复写入数据,需找到保留检查点同时避免重复的方法。
报错信息(翻译后)
com.databricks.sql.transaction.tahoe.DeltaUnsupportedOperationException:
[DELTA_SOURCE_TABLE_IGNORE_CHANGES] 在版本8的源表中检测到数据更新(例如CREATE OR REPLACE TABLE AS SELECT操作,参数Map(partitionBy -> [], clusterBy -> [], description -> null, isManaged -> true, properties -> {"delta.enableDeletionVectors":"true"}, statsOnLoad -> false))。当前不支持这种操作。如果这类操作会定期发生且你可以接受跳过变更,请设置选项'skipChangeCommits'为'true'。如果你希望数据更新能被反映,请使用新的检查点目录重启查询,或者如果使用DLT则执行全量刷新。如果需要处理这类变更,请切换到物化视图(MV)。源表路径为gs://databricks.....
原因分析
Delta流读取的检查点会记录源表的版本链信息,当源表执行重写类操作(如CREATE OR REPLACE TABLE、DROP PARTITION后重建等)时,源表的版本链会断裂,流任务无法从旧检查点记录的版本继续追踪增量,因此触发报错。
可行解决方案
1. 跳过不支持的变更(保留检查点)
如果可以接受丢失本次重写操作的数据,在readStream中添加skipChangeCommits选项,让流任务跳过这类不支持的变更,继续从后续版本追踪增量,无需重置检查点:
修改后的代码示例:
streaming_transactions = spark.readStream.format("delta") \ .option("cloudFiles.format", "deltaSharing") \ .option("skipChangeCommits", "true") # 新增选项:跳过不支持的变更 .table(f"{source_root_path}.{table_name}") \ .selectExpr("*", *metadata) streaming_transactions.writeStream.format("delta") \ .partitionBy("retrieved_datetime") \ .trigger(availableNow=True) \ .option("checkpointLocation", checkpoint) \ .option("readChangeFeed", "true") \ .option("mergeSchema", "true") \ .toTable( tableName=target_table_name, format="delta", outputMode="append", path=target_path )
2. 改用物化视图(Materialized Views)
如果需要完整捕获源表的重写类变更,官方推荐使用物化视图。MV会自动处理源表的结构变更、重写操作,无需手动维护检查点,同时能保证数据的一致性,避免重复写入。
3. 规范源表操作
如果源表的重写操作是可控的,尽量使用**增量写入(append)**代替CREATE OR REPLACE TABLE这类重写操作,保持源表版本链的连续性,这样流任务可以正常通过检查点追踪增量,不会触发报错。
内容的提问来源于stack exchange,提问作者Diego

