【发布时间】:2022-11-03 23:51:05
【问题描述】:
部署启用了自动缩放的数据流流作业时,它使用单个工作器。 假设管道读取 pubsub 消息,执行一些 DoFn 操作并上传到 BQ。 我们还假设 PubSub 队列已经有点大了。 所以管道开始并加载一些 pubsubs 在单个工作人员上处理它们。 几分钟后,它意识到需要一些额外的工人并创建它们。 许多 pubsub 消息已加载并正在处理但尚未确认。 这是我的问题:数据流将如何管理那些尚未确认的、正在处理的元素?
我的观察表明,数据流将许多已经处理的消息发送给新创建的工作人员,我们可以看到两个工作人员同时处理相同的元素。 这是预期的行为吗?
另一个问题是——下一步是什么?首胜?还是新的胜利? 我的意思是,我们有相同的 pubsub 消息,它仍在第一个工作人员和新工作人员上处理。 如果第一个工作人员的处理速度更快并完成处理怎么办?它将被确认并进入下游或将被丢弃,因为该元素的新进程已启动并且只有新进程才能完成?
【问题讨论】:
标签: google-cloud-dataflow apache-beam autoscaling