【问题标题】:Why do my Dask Futures get stuck in 'pending' and never finish?为什么我的 Dask Futures 卡在“待处理”中并且永远不会完成?
【发布时间】:2020-09-21 15:30:51
【问题描述】:

我有一些长时间运行的代码(大约 5 到 10 分钟的处理),我正试图以 Dask Future 的形式运行。这是一系列的几个离散步骤,我可以作为一个函数运行:

result : Future = client.submit(my_function, arg1, arg2)

或者我可以分成中间步骤:

# compose the result from the same intermediate results but with Futures
intermediate1 = client.submit(my_function1, arg1)
intermediate2 = client.submit(my_function2, arg1, arg2)
intermediate3 = client.submit(my_function3, intermediate2, arg1)
result = client.submit(my_function4, intermediate3)

如果我在本地运行它(例如,result = my_function(arg1, arg2)),它就会完成。如果我将它提交给 Dask,我会立即收到我的 Future - 正如预期的那样 - 但这项工作永远不会完成。此外,如果我将result.key 用作跟踪作业状态的一种方式,然后将未来重构为result = Future(key),它总是具有pending 的状态。

我想首先让它按原样运行,这样我就可以将我的处理工作卸载给我的 Dask 工作人员,而不是处理请求的 API,然后我希望能够开始跨节点拆分工作,所以我可以提高性能。但为什么我的工作就这么消失了?查看我的 Dask 调度程序 Web 界面,作业似乎都没有出现。但我知道 Dask 正在工作,因为我可以从我的 Jupyter 笔记本向它提交代码。

我正在从 Flask 服务器调用 client.submit,并且我正在返回密钥以便以后可以使用它。大致:

@app.route('/submit')
def submit():
    # ...
    future = client.submit(my_function, arg1, arg2)
    return jsonify({"key": future.key})

@app.route('/status/<key>')
def status(key):
    future = Future(key)
    return jsonify({"status": future.status})

当我的应用程序部署到 Kubernetes 时,我的 /submit 路由会返回一个 Future 密钥,但我的 Dask 状态页面没有显示任何处理任务。如果我在本地运行 Flask,我确实会看到一个任务出现,并且我的作业的输出确实在预期的延迟后出现;但是,当我使用从/submit 返回的 Future 键点击我自己的 /status/&lt;key&gt; 路径时,它总是显示状态为 pending

【问题讨论】:

    标签: python dask dask-distributed dask-kubernetes


    【解决方案1】:

    如果指向某个任务的所有未来都消失了,那么 Dask 可以随意忘记该任务。这允许 Dask 清理工作,而不是让所有中间结果永远存在。

    如果您想保留引用,那么您需要保留期货。这告诉 Dask 你仍然关心结果。您可以通过创建字典在您的烧瓶应用中本地执行此操作。

    futures = {}
    
    @app.route('/submit')
    def submit():
        # ...
        future = client.submit(my_function, arg1, arg2)
        futures[future.key] = future
        return jsonify({"key": future.key})
    
    @app.route('/status/<key>')
    def status(key):
        future = futures[key]
        return jsonify({"status": future.status})
    

    但您还需要考虑何时可以清理和释放这些期货。用这种方法,你会慢慢填满你的记忆。

    【讨论】:

    • 事实上,这正是我昨天所做的,但没有回来更新我的问题。
    • 有没有比将它存储在我们的 Flask 应用程序中更好的方法呢?我们不能告诉 Dask 以某种方式坚持未来吗?
    • 由于 Flask 通常与一些 WSGI 服务器一起使用,请求通常可以由不同的进程/线程处理,所以这不会冒给定未来不在字典中的风险吗?或者更一般地说,有没有办法在多进程环境中完全执行此示例中所做的操作?我一直在尝试使用Client.datasets 让集群持有期货,但也遇到了status 中的future 永远是'pending' 的问题。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-09-15
    • 2018-09-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多