【问题标题】:Distributing graphs to across cluster nodes将图分布到集群节点
【发布时间】:2018-02-11 21:20:44
【问题描述】:

我在 Dask.delayed 方面取得了不错的进展。作为一个团队,我们决定将更多时间用于使用 Dask 处理图表。

我有一个关于分发的问题。我在我们的集群上看到以下行为。我启动例如8 个节点上的每个节点都有 8 个工作人员,每个节点有 4 个线程,比如说/我然后 client.compute 8 个图表来创建模拟数据以供后续处理。我想让 8 个数据集为每个节点生成一个。但是,似乎发生的情况并非没有道理,这八个功能在前两个节点上运行。随后的计算在第一和第二节点上运行。因此,我看到缺乏缩放。随着时间的推移,其他节点将从诊断工作者页面中消失。这是预期的吗?

所以我想先按节点分配数据创建功能。所以当我想计算图表时,我现在这样做:

if nodes is not None:
    print("Computing graph_list on the following nodes: %s" % nodes)
    return client.compute(graph_list, sync=True, workers=nodes, **kwargs)
else:
    return client.compute(graph_list, sync=True, **kwargs)

这似乎设置正确:诊断进度条显示我的数据创建功能在内存中,但它们没有启动。如果省略节点,则计算按预期进行。此行为同时出现在集群和我的桌面上。

更多信息:查看调度程序日志,我确实看到了通信故障。

more dask-ssh_2017-09-04_09\:52\:09/dask_scheduler_sand-6-70\:8786.log
distributed.scheduler - INFO - -----------------------------------------------
distributed.scheduler - INFO -   Scheduler at:    tcp://10.143.6.70:8786
distributed.scheduler - INFO -       bokeh at:              0.0.0.0:8787
distributed.scheduler - INFO -        http at:              0.0.0.0:9786
distributed.scheduler - INFO - Local Directory:    /tmp/scheduler-ny4ev7qh
distributed.scheduler - INFO - -----------------------------------------------
distributed.scheduler - INFO - Register tcp://10.143.6.73:36810
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.143.6.73:36810
distributed.scheduler - INFO - Register tcp://10.143.6.71:46656
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.143.6.71:46656
distributed.scheduler - INFO - Register tcp://10.143.7.66:42162
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.143.7.66:42162
distributed.scheduler - INFO - Register tcp://10.143.7.65:35114
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.143.7.65:35114
distributed.scheduler - INFO - Register tcp://10.143.6.70:43208
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.143.6.70:43208
distributed.scheduler - INFO - Register tcp://10.143.7.67:45228
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.143.7.67:45228
distributed.scheduler - INFO - Register tcp://10.143.6.72:36100
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.143.6.72:36100
distributed.scheduler - INFO - Register tcp://10.143.7.68:41915
distributed.scheduler - INFO - Starting worker compute stream, tcp://10.143.7.68:41915
distributed.scheduler - INFO - Receive client connection: 5d1dab2a-914e-11e7-8bd1-180373ff6d8b
distributed.scheduler - INFO - Worker 'tcp://10.143.6.71:46656' failed from closed comm: Stream is clos
ed
distributed.scheduler - INFO - Remove worker tcp://10.143.6.71:46656
distributed.scheduler - INFO - Removed worker tcp://10.143.6.71:46656
distributed.scheduler - INFO - Worker 'tcp://10.143.6.73:36810' failed from closed comm: Stream is clos
ed
distributed.scheduler - INFO - Remove worker tcp://10.143.6.73:36810
distributed.scheduler - INFO - Removed worker tcp://10.143.6.73:36810
distributed.scheduler - INFO - Worker 'tcp://10.143.6.72:36100' failed from closed comm: Stream is clos
ed
distributed.scheduler - INFO - Remove worker tcp://10.143.6.72:36100
distributed.scheduler - INFO - Removed worker tcp://10.143.6.72:36100
distributed.scheduler - INFO - Worker 'tcp://10.143.7.67:45228' failed from closed comm: Stream is clos
ed
distributed.scheduler - INFO - Remove worker tcp://10.143.7.67:45228
distributed.scheduler - INFO - Removed worker tcp://10.143.7.67:45228
(arlenv) [hpccorn1@login-sand8 performance]$

这会引发任何可能的原因吗?

谢谢, 蒂姆

【问题讨论】:

    标签: distributed dask


    【解决方案1】:

    Dask 如何选择将任务分配给工作人员是很复杂的,并且会考虑负载平衡、数据传输、资源限制等许多问题。如果没有具体而简单的信息,就很难推断事情最终会走向何方例子。

    您可以尝试的一件事是一次提交所有计算,这可以让调度程序做出更明智的决策,而不是一次只看到一个。

    所以,你可以试试这样替换代码:

    futures = [client.compute(delayed_value) for delayed_value in L]
    wait(futures)
    

    这样的代码

    futures = client.compute(L)
    wait(futures)
    

    但老实说,我只给了 30% 的机会来解决您的问题。如果不深入研究您的问题,就很难知道发生了什么。如果您能提供一个非常简单可重现的代码示例,那将是最好的。

    【讨论】:

    • 为了让发行版正常工作,我必须将名称转换为 IP 地址。
    • 我的问题是内存订阅不足,因此即使将所有数据移动到一个节点(从 16 个节点),它仍然只有 25% 已满。所以我从 8 个节点/16 个线程和百分之几的内存开始,然后将它们缩小到一个工作中。我需要更多地考虑如何进行。
    • 我认为这与内存大小无关。查看调度程序日志。
    • 马特,我已经尽我所能。测试在单个节点(笔记本电脑、台式机和集群单节点)上运行正常。集群上的相同测试(即分布在 2、4、8、16 上)显示了相同的行为:节点下降到非常低的 CPU 并在一段时间后被删除。删除事件显示在诊断中,因此它看起来不像是错误。我现在的问题是如何控制这种行为。我可以告诉调度程序不要迁移吗?或类似的东西。谢谢,蒂姆
    • 恐怕我对您的问题了解不多,无法提供帮助。例如,我不知道您告诉调度程序不要迁移是什么意思。您可能想要创建一个mcve,其他人可以尝试轻松重现您观察到的问题。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-06-03
    • 2021-11-10
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多