【问题标题】:nested dask.compute not blocking嵌套的 dask.compute 不阻塞
【发布时间】:2017-08-18 01:54:23
【问题描述】:

dask.compute(...) 应该是一个阻塞调用。但是,当我嵌套了 dask.compute,而内部的执行 I/O(如 dask.dataframe.read_parquet)时,内部的 dask.compute 不会阻塞。这是一个伪代码示例:

import dask, distributed

def outer_func(name):
    files = find_files_for_name(name)
    df = inner_func(files).compute()
    # do work with df
    return result

def inner_func(files):
    tasks = [ dask.dataframe.read_parquet(f) for f in files ]
    tasks = dask.dataframe.concat(tasks)
    return tasks

client = distributed.Client(scheduler_file=...)
results = dask.compute([ dask.delay(outer_func)(name) for name in names ])

如果我启动 2 个工人,每个工人有 8 个进程,例如:

dask-worker --scheduler-file $sched_file --nprocs 8 --nthreads 1

,那么我预计最多 2 x 8 个并发的 inner_func 正在运行,因为 inner_func(files).compute() 应该是阻塞的。然而,我观察到的是,在一个工作进程中,一旦它开始 read_parquet 步骤,可能会有另一个 inner_func(files).compute() 开始运行。所以最后可能有多个inner_func(files).compute()在运行,有时会导致内存不足的错误。

这是预期的行为吗?如果是这样,是否有任何方法可以强制每个工作进程执行一个 inner_func(files).compute()?

【问题讨论】:

  • 这里似乎有点混乱。 dask.dataframe 创建惰性对象,在同样延迟/计算的函数中创建/计算这些对象是不正常的。考虑一下,这个函数被发送给一个工作者:你希望计算发生在哪里?
  • 这个例子中的嵌套是真实世界数据流恕我直言的典型。使用像 dask DataFrame 这样的分布式数据结构实际上并不总是可行/可取的,这样我们就可以避免这种嵌套。因为 dask DataFrame API 比 pandas 小,并且因为保持一个有效的串行代码版本非常重要。从我所看到的,inner_func 似乎在 dask-worker 进程内的多个线程中运行,但我只使用例如:dask-worker --scheduler-file sched.json --nprocs 3 --nthreads 1 为每个工作者指定一个线程--local-directory /tmp/

标签: python dask dask-distributed dask-delayed


【解决方案1】:

当您要求 dask 分布式调度程序运行工作时,它会将函数的代码和所需的任何数据发送到不同进程中的工作函数,可能在不同的机器上。这些工作进程忠实地执行函数,并像正常的 python 代码一样运行。关键是,正在运行的函数不知道它在 dask worker 上——默认情况下,它会看到没有设置全局 dask 分布式客户端,并执行 dask 通常在这种情况下会做的事情:执行任何 dask默认调度程序(线程调度程序)上的工作负载。

如果您确实必须在任务中执行完整的 dask-compute 操作,并且希望这些操作使用运行这些任务的分布式调度程序,您将需要使用 worker client。但是,我认为在您的情况下,改写作业以删除嵌套(类似于上面的伪代码,尽管这也可以与计算一起使用)可能是更简单的方法。

【讨论】:

    【解决方案2】:

    多进程调度器似乎不是这种情况。

    为了使用分布式调度程序,我通过distributed.Client API 使用有节奏的作业提交而不是依赖dask.compute 找到了解决方法。 dask.compute 可以用于简单的用例,但显然不知道可以安排多少未完成的任务,因此在这种情况下会超出系统。

    这是运行 dask.Delayed 任务集合的伪代码:

    import distributed as distr
    
    def paced_compute(tasks, batch_size, client):
        """
        Run delayed tasks, maintaining at most batch_size running at any
        time. After the first batch is submitted,
        submit a new job only after an existing one is finished, 
        continue until all tasks are computed and finished.
    
        tasks: collection of dask.Delayed
        client: distributed.Client obj
        """
        results, tasks = [], list(tasks)
        working_futs = client.compute(tasks[:batch_size])
        tasks = tasks[batch_size:]
        ac = distr.as_completed(working_futs)
        for fut in ac:
            res = fut.result()
            results.append(res)
            if tasks:
                job = tasks.pop()
                ac.add(client.compute(job))
        return results
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2010-10-31
      • 2013-10-15
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多