【问题标题】:How to structure a Dask application that processes a fixed number of inputs from a queue?如何构建处理来自队列的固定数量输入的 Dask 应用程序?
【发布时间】:2017-08-23 20:22:34
【问题描述】:

我们需要执行以下操作。给定一个将提供已知数量消息的 Redis 通道:

  1. 对于从通道消费的每条消息:

    • 从 Redis 获取 JSON 文档
    • 解析 JSON 文档,提取结果对象列表
  2. 聚合所有结果对象以生成单个结果

我们希望将第 1 步和第 2 步分配给许多工作人员,并避免将所有结果都收集到内存中。我们还想显示两个步骤的进度条。

但是,我们看不到一种构建应用程序的好方法,以便我们可以看到进度并在系统中不断进行工作,而不会因为不合时宜的时间而阻塞。

例如,在第 1 步中,如果我们从 Redis 通道读取到队列中,那么我们可以将队列传递给 Dask,在这种情况下,我们开始处理每条消息,而无需等待所有消息。但是,如果我们使用队列,我们​​就看不到显示进度的方法(大概是因为队列通常具有未知大小?)

如果我们从 Redis 通道收集到一个列表并将其传递给 Dask,那么我们可以看到进度,但是我们必须等待来自 Redis 的所有消息,然后才能开始处理第一个消息。

有解决此类问题的推荐方法吗?

【问题讨论】:

    标签: dask dask-distributed


    【解决方案1】:

    如果您的 Redis 通道是并发访问安全的,那么您可能会提交许多未来以从通道中提取元素。这些将在不同的机器上运行。

    from dask.distributed import Client, progress
    client = Client(...)
    
    futures = [client.submit(pull_from_redis_channel, ..., pure=False) for _ in range(n_items)]
    futures2 = client.map(process, futures)
    
    progress(futures2)
    

    【讨论】:

    • 非常好!我使用的是 Redis pub/sub,它对此不起作用,但是 LPUSH + BLPOP 就像一个点点队列,并且与上述方法非常配合,谢谢
    猜你喜欢
    • 2023-03-29
    • 2022-08-06
    • 1970-01-01
    • 2019-06-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多