【问题标题】:How does dataflow manage current processes during upscaling streaming job?在升级流作业期间,数据流如何管理当前进程?
【发布时间】:2022-11-03 23:51:05
【问题描述】:

部署启用了自动缩放的数据流流作业时,它使用单个工作器。 假设管道读取 pubsub 消息,执行一些 DoFn 操作并上传到 BQ。 我们还假设 PubSub 队列已经有点大了。 所以管道开始并加载一些 pubsubs 在单个工作人员上处理它们。 几分钟后,它意识到需要一些额外的工人并创建它们。 许多 pubsub 消息已加载并正在处理但尚未确认。 这是我的问题:数据流将如何管理那些尚未确认的、正在处理的元素?

我的观察表明,数据流将许多已经处理的消息发送给新创建的工作人员,我们可以看到两个工作人员同时处理相同的元素。 这是预期的行为吗?

另一个问题是——下一步是什么?首胜?还是新的胜利? 我的意思是,我们有相同的 pubsub 消息,它仍在第一个工作人员和新工作人员上处理。 如果第一个工作人员的处理速度更快并完成处理怎么办?它将被确认并进入下游或将被丢弃,因为该元素的新进程已启动并且只有新进程才能完成?

【问题讨论】:

    标签: google-cloud-dataflow apache-beam autoscaling


    【解决方案1】:

    Dataflow 提供对每条记录的一次性处理。有趣的是,这并不意味着用户代码每条记录只运行一次,无论是通过流式运行程序还是批处理运行程序。

    它可能通过用户转换多次运行给定记录,甚至可能在多个工作人员上同时运行相同的记录;这对于保证在面对工人故障时至少处理一次是必要的。这些调用中只有一个可以“获胜”并在管道的下游产生输出。

    更多信息在这里-https://cloud.google.com/blog/products/data-analytics/after-lambda-exactly-once-processing-in-google-cloud-dataflow-part-1

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2011-07-18
      • 1970-01-01
      • 2022-06-18
      • 2020-07-11
      • 2020-12-04
      相关资源
      最近更新 更多