astronomy-commons/lsdb

Specify memory required to merge computed worker results

開放

#678 建立於 2025年4月4日

 (1 則留言) (0 個反應) (0 位負責人)Python (21 個分叉)auto 404
enhancementhelp wanted

倉庫指標

星標
 (54 顆星)
PR 合併指標
 (PR 指標待抓取)

描述

When creating a Dask Client we specify the number of workers and the memory limit for each of them. This means that each worker is assigned the same amount of memory. E.g.:

dask.distributed.Client(n_workers=16, memory_limit="4GiB")

When we hit compute one of those workers is responsible for merging the results of all individual computations and, if the result is bigger than that single worker's memory, the computation fails. We should be able to specify a different amount of memory required for a final worker to complete the computation.

There seems to be no way of specifying which merge worker to use on compute of a Dask DataFrame but that seems to be possible through the Client: https://distributed.dask.org/en/stable/api.html#distributed.Client.persist. Requires some investigation.

貢獻者指南