【问题标题】:Dask job queue design pattern?Dask作业队列设计模式?
【发布时间】:2019-12-05 12:55:34
【问题描述】:

假设我有一个简单的高成本函数,可以将一些结果存储到文件中:

def costly_function(filename):
    time.sleep(10)
    with open('filename', 'w') as f:
        f.write("I am done!)

现在假设我想在 dask 中安排一些这样的任务,然后异步接收这些请求并一一运行这些函数。我目前正在设置一个 dask 客户端对象...

cluster = dask.distributed.LocalCluster(n_workers=1, processes=False)  # my attempt at sequential job processing
client = dask.distributed.Client(cluster)

...然后以交互方式(来自 IPython)调度这些作业:

>>> client.schedule(costly_function, "result1.txt")
>>> client.schedule(costly_function, "result2.txt")
>>> client.schedule(costly_function, "result3.txt")

我遇到的问题是这些任务不是连续运行而是并行运行,这在我的特定情况下会导致并发问题。

所以我的问题是:像我上面在 dask 中描述的那样设置作业队列的正确方法是什么?

【问题讨论】:

    标签: python dask job-queue


    【解决方案1】:

    好的,我想我可能有一个解决方案(不过,请随时提出更好的解决方案!)。它需要稍微修改之前的代价函数:

    def costly_function(filename, prev_job=None):
        time.sleep(10)
        with open('filename', 'w') as f:
            f.write("I am done!")
    
    cluster = dask.distributed.LocalCluster(n_workers=1, processes=False)  # my attempt at sequential job processing
    client = dask.distributed.Client(cluster)
    

    然后在交互式上下文中,您将编写以下内容:

    >>> future = client.submit(costly_function, "result1.txt")
    >>> future = client.submit(costly_function, "result2.txt", prev_job=future)
    >>> future = client.submit(costly_function, "result3.txt", prev_job=future)
    

    【讨论】:

    • 我稍微修改了你的答案。您无需致电.result。这是自动完成的。此外,方法名称是提交,而不是调度。
    • 嘿,感谢您的编辑!你能解释一下为什么在这种情况下不需要调用 .result() 吗?我不知道这究竟是如何自动完成的。
    • 当您将未来作为参数包含在提交调用中时,Dask 将其识别为数据依赖项。它会等到未来完成计算后再运行新任务,并传入计算结果,而不是未来。您可以在docs.dask.org/en/latest/futures.html 了解有关 Dask 期货的更多信息
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2022-08-18
    • 1970-01-01
    • 1970-01-01
    • 2012-02-09
    • 2015-07-07
    • 2018-06-15
    • 2011-05-17
    相关资源
    最近更新 更多