【问题标题】:How to set processing timeouts in apache beam / Dataflow python batch jobs?如何在 apache beam / Dataflow python 批处理作业中设置处理超时?
【发布时间】:2020-04-01 11:10:47
【问题描述】:

我目前正在使用 stopit 库 https://github.com/glenfant/stopit 来设置批处理作业中的每个元素处理超时。这些作业在直接运行器上运行,我可以超时执行耗时过长的函数。

为批处理作业设置每个元素处理超时的梁方式是什么?

有没有一种方法可以为数据流批处理作业设置处理超时?

我的用例是从文本中提取命名实体。如果正在处理的文档过长,NER 过程有时会花费很长时间。

摆脱这种依赖并转向Beam原生解决方案会很好。

【问题讨论】:

  • 根据 Apache Beam 文档,您可以将 waituntilfinish() 方法与 direct runnerdata flow runner 一起使用。使用此方法,您可以定义管道将超时的时间量。这就是你要找的吗?
  • 我正在寻找每个元素的进程超时而不是整个管道超时。我也更新了这个问题以澄清这一点。不过感谢您的有用建议!
  • 您能否举例说明您在做什么以及为什么要为每个元素设置超时?所以我可以更好地理解并帮助你。
  • 对每一段文本运行一个NLP提取管道,代码的执行时间取决于文本。对于特定文本,我使用的 NLP 库运行速度可能非常缓慢,我想终止运行缓慢的提取并在另一台机器上重新处理它们。

标签: python-3.x google-cloud-dataflow apache-beam


【解决方案1】:

据我了解,问题的答案是及时处理,也保持状态。 假设您有一个函数 f 并获取将使用该函数的任何批次的输出。 所以基本上我们必须在收到批次后标记批次(更新状态),我们将设置计时器,根据可以设置的水印/到期时间更新输出,如果有任何输出,我们将收到它并且如果我们没有根据您的查询从函数中输出任何输出,我们当然可以重新路由该批次。

这不是一个精确的解决方案,但可以解决

为了更好地理解你可以在这里阅读:apache beam documentation

【讨论】:

  • 有趣,感谢分享上次我检查我认为及时处理只支持流式作业。看起来它可以用于流式传输。我稍后会研究它。
猜你喜欢
  • 1970-01-01
  • 2021-01-19
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多