Apache Beam Python ReadFromPubsub IO内存泄漏问题排查求助
Dataflow流处理管道内存泄漏排查问题
问题背景
我们运行着一个Dataflow流处理管道,从Pub/Sub订阅读取消息,将字典转换为dataclass后写入PostgreSQL。近期发现Pub/Sub吞吐量偶尔会降至0,此时内存利用率通常达95%以上,但多数情况下高内存占用时数据仍能稳定流动。
为排查问题,我逐个从后往前移除PTransform并部署观察,问题在所有场景下都存在——甚至当管道仅包含单个ReadFromPubsub转换时也是如此,因此怀疑是库实现存在内存泄漏。
观察与疑问
- 内存占用下降通常对应扩缩容事件:从1个Worker扩到2个时,原Worker A有时会被替换为Worker B和C,是否意味着Dataflow检测到Worker A内存问题并将其关闭?其他时候内存下降对应Worker A重启,这也会清除内存。
- 论坛中有回复建议用窗口作为解决方案,我的理解是:在仅使用全局窗口的流管道中(我的场景),DoFns一直处于运行状态,没有垃圾回收机会(
teardown从未被调用?),因此需要通过创建窗口来强制触发,比如设5秒窗口、有1个元素就触发,再按键分组并展平,本质是无操作仅为触发窗口机制。
环境配置
- Apache Beam版本:
2.48.0 - Python版本:
3.9.14 - 运行器:
Dataflow
内容的提问来源于stack exchange,提问作者leech
相关产品推荐
相关产品推荐

