# @distributed fails for many workers?

**URL:** <https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145>\
**Category:** Julia at Scale\
**Created:** [May 12, 2019, 5:33pm UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145 "2019-05-12T17:33:50Z")\
**Posts on this page:** 20\
**Page:** 1

<div class="post-metadata">

**Author:** ![dkiese](https://avatars.discourse-cdn.com/v4/letter/d/71e660/32.png) [@dkiese](https://discourse.julialang.org/u/dkiese)\
**Post date:** [May 12, 2019, 5:33pm UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/1 "2019-05-12T17:33:50Z")

</div>

Hello folks,

I am currently running some data analysis on a cluster with 68 cores per node. To this end I set up a SharedArray for 68 workers and let them operate on it via @sync @distributed. For small datasets this works fine on my local machine.

If I however launch my Julia script on the cluster node, the workers seem to connect to the master (at least nprocs() = 69) before the iteration, but then I see at most 3 workers really doing something using top. Simultaneously the master keeps consuming more and more memory until the code crashes. I really do not get the problem behind this, since locally I cannot observe any memory leakage. Any ideas how to test what’s going on with the workers?

I am using Julia 1.0.1.

Hoping somebody can help out.

EDIT 1: Below you can find a MWE that causes this for me.

```julia
using Distributed
using SharedArrays

addprocs(68, topology=:master_worker)

A = ones(Float64, 100000000)
B = SharedArray{Float64, 1}((length(A)))

@sync @distributed for i in 1 : length(A)
   B[i] = A[i]
end

```

EDIT 2:

The MWE also works without the shared array. Use e.g. `println(A[i])` in the loop. This yields

```julia
From Worker N: 1.0

```

But the workers print their id (=N) sequentially, only one operating at a time. So first N = 2, then N = 7, then some other id and so on.

EDIT 3:

After further investigation, it seems to me that I am doing something wrong when allocating resources for an interactive session or when submitting the job. Namely, I guess all processes are started on the same core, causing them to execute serially and not in parallel. Can somebody explain how to allocate resources on a slurm cluster for the above shared memory problem? So far I used

`salloc --nodes=1 --ntasks=1 --cpus-per-task=68`

and then started the julia code from above via `julia MWE.jl`.

---

<div class="post-metadata">

**Author:** ![jpsamaroo](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/jpsamaroo/32/46804_2.png) [@jpsamaroo](https://discourse.julialang.org/u/jpsamaroo)\
**Post date:** [May 13, 2019, 12:15pm UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/2 "2019-05-13T12:15:51Z")

</div>

Do the allocations of `A` and `B` succeed without crashing? And can you confirm that all 68 workers remain alive while that loop is running (they don’t get killed by OOM or something similar)?

---

<div class="post-metadata">

**Author:** ![dkiese](https://avatars.discourse-cdn.com/v4/letter/d/71e660/32.png) [@dkiese](https://discourse.julialang.org/u/dkiese)\
**Post date:** [May 13, 2019, 12:20pm UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/3 "2019-05-13T12:20:14Z")

</div>

The allocation succeeds. procs(B) shows all 68 workers, so B should also be properly mapped. How would you confirm the latter? From top I only see that some small number of workers (2 or 3) is active with the rest apparently sleeping. The IDs of the respective processes change though.

---

<div class="post-metadata">

**Author:** ![jpsamaroo](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/jpsamaroo/32/46804_2.png) [@jpsamaroo](https://discourse.julialang.org/u/jpsamaroo)\
**Post date:** [May 13, 2019, 7:14pm UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/5 "2019-05-13T19:14:44Z")

</div>

Are those IDs actual workers, or just other threads of the master process? It’s possible the libuv threads of the master process are doing all that work for whatever reason, which wouldn’t be very obvious from the output of top (which is why I use `htop`).

---

<div class="post-metadata">

**Author:** ![johnh](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/johnh/32/3615_2.png) [@johnh](https://discourse.julialang.org/u/johnh)\
**Post date:** [May 13, 2019, 7:21pm UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/6 "2019-05-13T19:21:56Z")

</div>

These are separate processes, not threads.  
I see many processes sitting in epoll\_pwait as three or four processes are running.

If I run thei minimal example I do not get a crash - ARM 64 platform, lots of RAM, Julia 1.1

---

<div class="post-metadata">

**Author:** ![jpsamaroo](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/jpsamaroo/32/46804_2.png) [@jpsamaroo](https://discourse.julialang.org/u/jpsamaroo)\
**Post date:** [May 13, 2019, 7:41pm UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/7 "2019-05-13T19:41:33Z")

</div>

Julia itself uses threading for each process (via libuv) to service things like syscalls and other blocking operations. My point was that, if for some reason you were hitting an endless stream of syscalls, it would probably look like 2-3 threads running eternally (although I think libuv starts 4 by default). Especially if you’re seeing a ton of `epoll_pwait`, which is what libuv calls to wait on its set of file descriptors.

---

<div class="post-metadata">

**Author:** ![dkiese](https://avatars.discourse-cdn.com/v4/letter/d/71e660/32.png) [@dkiese](https://discourse.julialang.org/u/dkiese)\
**Post date:** [May 13, 2019, 7:57pm UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/8 "2019-05-13T19:57:32Z")

</div>

I am actually not sure if they are threads or processes, I’d have to check. How do u tell from htop?

---

<div class="post-metadata">

**Author:** ![dkiese](https://avatars.discourse-cdn.com/v4/letter/d/71e660/32.png) [@dkiese](https://discourse.julialang.org/u/dkiese)\
**Post date:** [May 13, 2019, 7:58pm UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/9 "2019-05-13T19:58:50Z")

</div>

But you can reproduce that not all workers participate equally in executing the loop?

---

<div class="post-metadata">

**Author:** ![johnh](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/johnh/32/3615_2.png) [@johnh](https://discourse.julialang.org/u/johnh)\
**Post date:** [May 13, 2019, 8:01pm UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/10 "2019-05-13T20:01:07Z")

</div>

Remember not all workers in a parallel computation are guaranteed to finish at the same time.

---

<div class="post-metadata">

**Author:** ![jpsamaroo](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/jpsamaroo/32/46804_2.png) [@jpsamaroo](https://discourse.julialang.org/u/jpsamaroo)\
**Post date:** [May 13, 2019, 8:20pm UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/11 "2019-05-13T20:20:56Z")

</div>

Once you enable “Tree view” in the `htop` settings, `htop` shows them branching off of the tree just like processes, but puts them (in my case) in a slightly dimmer color than child processes:  
 ![htop](https://global.discourse-cdn.com/julialang/original/3X/8/8/88656e085419f3f881c641052fa7ef2337988f83.png)

In the above image, `nvim` has one thread (also called `nvim`), and also one child process (called `languageclient`).

---

<div class="post-metadata">

**Author:** ![dkiese](https://avatars.discourse-cdn.com/v4/letter/d/71e660/32.png) [@dkiese](https://discourse.julialang.org/u/dkiese)\
**Post date:** [May 14, 2019, 4:21am UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/12 "2019-05-14T04:21:55Z")

</div>

Sure. But they should at least all start more or less simultaneously and not one after another

---

<div class="post-metadata">

**Author:** ![dkiese](https://avatars.discourse-cdn.com/v4/letter/d/71e660/32.png) [@dkiese](https://discourse.julialang.org/u/dkiese)\
**Post date:** [May 14, 2019, 4:22am UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/13 "2019-05-14T04:22:13Z")

</div>

Thanks. I’ll try that out

---

<div class="post-metadata">

**Author:** ![affans](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/affans/32/11911_2.png) [@affans](https://discourse.julialang.org/u/affans)\
**Post date:** [May 14, 2019, 4:46am UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/14 "2019-05-14T04:46:49Z")

</div>

I don’t know if your Edit has been responded to, but try running `srun` instead of `salloc`.

---

<div class="post-metadata">

**Author:** ![dkiese](https://avatars.discourse-cdn.com/v4/letter/d/71e660/32.png) [@dkiese](https://discourse.julialang.org/u/dkiese)\
**Post date:** [May 14, 2019, 5:20am UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/15 "2019-05-14T05:20:20Z")

</div>

It has not been responded to yet. Should have written this more explicitly however. After the salloc I do srun - - pty bash and start the script afterwards. Alternatively I can also call it directly from the Julia REPL. That should be the same as running directly with srun right?

---

<div class="post-metadata">

**Author:** ![dkiese](https://avatars.discourse-cdn.com/v4/letter/d/71e660/32.png) [@dkiese](https://discourse.julialang.org/u/dkiese)\
**Post date:** [May 14, 2019, 6:53am UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/16 "2019-05-14T06:53:09Z")

</div>

Tested it, these are indeed julia processes not process threads.

---

<div class="post-metadata">

**Author:** ![johnh](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/johnh/32/3615_2.png) [@johnh](https://discourse.julialang.org/u/johnh)\
**Post date:** [May 14, 2019, 9:07am UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/17 "2019-05-14T09:07:07Z")

</div>

You have a 68 core compute node? May I Ask what architecture the processors are? Intel, AMD or ARM?

---

<div class="post-metadata">

**Author:** ![dkiese](https://avatars.discourse-cdn.com/v4/letter/d/71e660/32.png) [@dkiese](https://discourse.julialang.org/u/dkiese)\
**Post date:** [May 14, 2019, 9:08am UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/18 "2019-05-14T09:08:37Z")

</div>

Its an Intel Architecture with KNL processors

---

<div class="post-metadata">

**Author:** ![johnh](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/johnh/32/3615_2.png) [@johnh](https://discourse.julialang.org/u/johnh)\
**Post date:** [May 14, 2019, 9:58am UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/19 "2019-05-14T09:58:08Z")

</div>

What happens when you use 48 workers?  
I ask since there is this in lubuv

```julia
assert(timeout >= -1);
  base = loop->time;
  count = 48; /* Benchmarks suggest this gives the best throughput. */
  real_timeout = timeout;

  for (;;) {
    /* See the comment for max_safe_timeout for an explanation of why
     * this is necessary. Executive summary: kernel bug workaround.
     */
    if (sizeof(int32_t) == sizeof(long) && timeout >= max_safe_timeout)
      timeout = max_safe_timeout;

    nfds = epoll_pwait(loop->backend_fd,
                       events,
                       ARRAY_SIZE(events),
                       timeout,
                       psigset);

```

---

<div class="post-metadata">

**Author:** ![johnh](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/johnh/32/3615_2.png) [@johnh](https://discourse.julialang.org/u/johnh)\
**Post date:** [May 14, 2019, 9:58am UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/20 "2019-05-14T09:58:29Z")

</div>

Someone please remind me how to quote code on here…

---

<div class="post-metadata">

**Author:** ![dkiese](https://avatars.discourse-cdn.com/v4/letter/d/71e660/32.png) [@dkiese](https://discourse.julialang.org/u/dkiese)\
**Post date:** [May 14, 2019, 9:58am UTC](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145/21 "2019-05-14T09:58:57Z")

</div>

use ``` before and after your code snippet

[Next page](https://discourse.julialang.org/t/distributed-fails-for-many-workers/24145.md?page=2)
