Hi everyone,
I’m implementing a distributed pipeline on a shared cluster where each job is computationally expensive and may run for several hours.
Because workers are allocated by a scheduler, they may become available gradually. I would like to start computing as soon as the first worker is ready, then let additional workers begin processing queued jobs as they join the pool.
The following local example simulates this situation:
using Distributed
pids = addprocs(2)
pool = CachingPool([pids[1]])
@async begin
sleep(1)
push!(pool, pids[2])
@info "Added second worker" pid=pids[2]
end
started = time()
results = pmap(pool, 1:3) do job
@info "Starting" job pid=myid() elapsed=time() - started
sleep(10)
return myid()
end
@show results
rmprocs(pids)
I expected the second worker to start one of the queued jobs shortly after it was added to pool, approximately one second after the beginning of the run. Instead, it appears that the new worker may remain idle until the first job finishes.
The number of scheduling tasks is checked while input elements are submitted. If submission is already blocked because all existing scheduling tasks are busy, increasing the pool size does not seem to create another scheduling task immediately. The new worker is only noticed after an existing job finishes and submission advances.
This raises a few questions:
- Is adding workers to a
WorkerPoolorCachingPoolwhilepmapis running an officially supported use case? - Is the delayed utilization of newly added workers expected, or could it be considered a scheduling limitation?
- What is the recommended pattern for an elastic worker pool where workers become available gradually?
- Would manually combining
asyncmapwithremotecall_fetch(f, pool, args...)be appropriate, so that tasks wait directly for workers to enter the pool? - Is there an existing way to tell
pmapthe expected or maximum number of workers independently of the pool’s current size?
In the actual application, addprocs runs asynchronously through a cluster manager, and every worker is initialized before being added to the computation pool. The example above uses already-created local workers only to reproduce the scheduling behavior.
Thanks!