Databricks环境PySpark Streaming使用Checkpoint出现OOM及数据丢失问题咨询
问题根因分析
- 内存泄漏问题:你使用的PySpark 3.0.1版本属于较老的发行版,存在Delta Streaming checkpoint元数据加载后未及时回收的已知缺陷。流任务长期运行过程中,checkpoint目录下累积的偏移量记录、批次提交日志会持续驻留在Driver内存中,运行约10天的累积量会超过Driver内存阈值触发OOM。重建checkpoint会清空所有历史元数据,所以内存泄漏问题会临时消失。
- 数据漏处理问题:你在写入流配置中开启了
ignoreChanges=true参数,该参数会跳过上游Delta表的变更一致性校验,当上游Delta表存在数据更新、删除操作时,流任务会直接忽略对应变更分区的处理,最终导致部分数据丢失。重建checkpoint后任务会重新扫描上游Delta的全量偏移量,所以会临时修复漏处理问题,但不调整参数的情况下后续仍会复现。 - 潜在代码缺陷:你当前代码中的过滤条件写法存在逻辑错误:
pd.Timestamp.utcnow() - pd.Timedelta(hours=1)是在流任务启动时(Driver端)一次性计算的固定时间值,后续所有运行批次都会用这个固定值做过滤,不会动态每批取最新的一小时时间。如果要实现每批动态过滤最近一小时的数据,需要替换为Spark原生的时间函数:F.col("TIME") > F.current_timestamp() - F.expr("INTERVAL 1 HOUR")。
补充问题解答
修改for_each_batch函数内容或者过滤条件是否需要重建checkpoint,分两种情况判断:
- 仅修改
for_each_batch内部业务逻辑,没有调整输入源、输出路径、流触发规则等配置时,不需要重建checkpoint。checkpoint仅存储流任务的偏移量、状态数据,不存储用户自定义函数的代码,重启任务后会自动加载新的函数逻辑处理后续批次。 - 修改输入流的过滤条件时,必须重建checkpoint。过滤条件直接作用于输入流的偏移量计算逻辑,原有checkpoint中存储的是基于旧过滤规则生成的偏移量,复用旧checkpoint会直接导致数据漏读或者重复读取。
内容的提问来源于stack exchange,提问作者Noé Achache
相关产品推荐
相关产品推荐

