【问题标题】:Repeated task execution using the distributed Dask scheduler使用分布式 Dask 调度程序重复执行任务
【发布时间】:2017-02-01 03:19:23
【问题描述】:

我正在使用 Dask 分布式调度程序,在本地运行一个调度程序和 5 个工作人员。我将delayed() 任务列表提交给compute()

当任务的数量是 20(一个数字 >> 比工人的数量)并且每个任务至少需要 15 秒时,调度程序开始重新运行一些任务(或者并行执行它们更多不止一次)

这是一个问题,因为任务修改了 SQL 数据库,如果它们再次运行,它们最终会引发异常(由于数据库唯一性约束)。我没有在任何地方设置pure=True(我相信默认是False)。除此之外,Dask 图很简单(任务之间没有依赖关系)。

仍然不确定这是 Dask 中的功能还是错误。我有一种直觉,这可能与工人偷窃有关......

【问题讨论】:

    标签: python dask


    【解决方案1】:

    正确,如果一项任务被分配给一个工作人员,而另一个工作人员变得空闲,它可能会选择从其他工作人员那里窃取多余的任务。它有可能会窃取刚刚开始运行的任务,在这种情况下,该任务将运行两次。

    处理此问题的干净方法是确保您的任务是幂等的,即使运行两次它们也会返回相同的结果。这可能意味着在您的任务中处理您的数据库错误。

    这是非常适合数据密集型计算工作负载但不适合数据工程工作负载的策略之一。设计一个同时满足这两种需求的系统是很棘手的。

    【讨论】:

    • 谢谢。至少现在我可以停止尝试调试它了。所以它的特点。也许您可以更新在线文档以强调这是一种可能性?另外,如果一个任务可以以任何一种方式运行多次,那么不确定使用“纯”参数有什么意义。
    • 我还提出了一个 github 问题,要求提供一种方法来关闭应该只运行一次的任务的工作窃取:github.com/dask/distributed/issues/847
    • 如果发生这种情况并且任务运行两次(例如),Dask 的 DAG 中使用哪个任务结果?这只是一场直线上升的比赛吗?
    • 这不是一场比赛,而是一个随机的选择。
    猜你喜欢
    • 1970-01-01
    • 2017-02-13
    • 2018-08-08
    • 1970-01-01
    • 2019-08-26
    • 2021-07-20
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多