如何获取Delta表最新插入提交时间,避免Spark Structured Streaming读取报错
解决方案
要规避startingTimestamp晚于Delta表最新提交时间的报错,你可以通过Delta Lake内置的历史记录接口提前查询表的最新提交时间,和你配置的起始时间做校验调整后再传入流读参数,具体实现如下:
实现步骤
- 第一步:导入DeltaTable工具类
from delta.tables import DeltaTable
- 第二步:获取目标Delta表的实例
如果你用的是元数据中注册的表名(和你示例代码中table(t)的入参一致),用如下方式获取:
delta_table = DeltaTable.forName(spark, t) # 如果是按路径存储的Delta表,替换为 DeltaTable.forPath(spark, "delta表的存储路径")
- 第三步:查询Delta表的最新提交时间
你可以根据需求取全局最新提交时间,或者过滤只取插入操作的最新提交时间:
# 方式1:取全局最新的提交时间(不管操作类型,推荐用这个,避免遗漏非插入的变更提交) latest_commit_ts = delta_table.history(1).collect()[0]["timestamp"] # 方式2:仅取插入/追加操作的最新提交时间 # latest_insert_ts = delta_table.history()\ # .filter("operation = 'WRITE' AND operationParameters.mode = 'Append'")\ # .orderBy("version", ascending=False)\ # .limit(1).collect()[0]["timestamp"]
- 第四步:校验并调整
startingTimestamp参数
import datetime # 注意如果你的starting_time_stamp是字符串格式,需要先转为datetime类型再比较 if isinstance(starting_time_stamp, str): starting_time_stamp = datetime.datetime.fromisoformat(starting_time_stamp) # 若配置的起始时间晚于表最新提交时间,直接替换为表最新提交时间,也可以根据业务需求抛出友好提示 if starting_time_stamp > latest_commit_ts: starting_time_stamp = latest_commit_ts
- 第五步:将校验后的参数传入原流读逻辑即可
df = ( spark.readStream.format("delta") .option("startingTimestamp", starting_time_stamp.isoformat()) .table(t) )
注意事项
history(1)方法仅返回最近1个版本的提交记录,性能开销极低,不会影响作业运行效率- 返回的
timestamp字段为标准datetime类型,和自定义时间比较时注意保持类型一致
内容的提问来源于stack exchange,提问作者MetallicPriest
相关产品推荐
相关产品推荐

