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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 10:36:03