# The ultimate guide to distributed computing

**URL:** <https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867>\
**Category:** Julia at Scale\
**Tags:** parallel, cluster, distributed\
**Created:** [June 22, 2020, 2:52pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867 "2020-06-22T14:52:10Z")\
**Posts on this page:** 20\
**Page:** 1

<div class="post-metadata">

**Author:** ![juliohm](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/juliohm/32/215266_2.png) [@juliohm](https://discourse.julialang.org/u/juliohm)\
**Post date:** [June 22, 2020, 2:52pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/1 "2020-06-22T14:52:10Z")

</div>

# HOWTO: Distributed computing

In this thread I would like to put together relevant information for people interested in distributing simple function calls on multiple cluster nodes. Although the task is simple, there are some rough (undocumented) corners in the language that inhibit even experienced users from accomplishing it currently.

The idea here is to update this content every now and then to reflect the latest (and cleanest) way of performing distributed computing with remote workers in Julia. If you read Discourse, you will find many related threads where people shared solutions for specific problems, which are currently outdated. I think we need a central thread of discussion to solve most issues once and for all.

## Sample script

We will consider a sample script that processes a set of files in a `data` folder and saves the results in a `results` folder. I like this task because it involves IO and file paths, which can get tricky in remote machines:

```julia-auto
# instantiate environment
using Pkg; Pkg.instantiate()

# load dependencies
using CSV

# HELPER FUNCTIONS
# ----------------
function process(infile, outfile)
  # read file from disk
  csv = CSV.read(infile)

  # perform calculations
  sleep(60) # pretend it takes time
  csv.new = rand(size(csv,1))

  # save new file to disk
  CSV.write(outfile, csv)
end

# MAIN SCRIPT
# -----------

# relevant directories
indir = "data"
outdir = "results"

# files to process
infiles = readdir(indir, join=true)
outfiles = joinpath.(outdir, basename.(infiles))
nfiles = length(infiles)

for i in 1:nfiles
  process(infiles[i], outfiles[i])
end

```

We follow Julia’s best practices:

1. We start by instantiating the environment in the host machine, which lives in the files `Project.toml` and `Manifest.toml` in the project directory (the same directory of the script).
2. We then load the dependencies of the project, and define helper functions to be used.
3. The main work is done in a loop that calls the helper function with various files.

Let’s call this script `main.jl`. We can `cd` into the project directory and call the script as follows (assuming Julia v1.4 or higher):

```shell
$ julia --project main.jl

```

* * *

## Parallelization (same machine)

Our goal is to process the files in parallel. First, we will make minor modifications to the script to be able to run it with multiple processes on the same machine (e.g. the login node). This step is important for debugging:

- We load the `Distributed` stdlib to replace the simple for loop by a `pmap` call. It seems that `Distributed` is always available so we don’t need to instantiate the environment before loading it. That will be important because we will instantiate the other dependencies in all workers with a `@everywhere` block call that is already available without any previous instantiation.
- We add worker processes with `addprocs` on the same machine and tell these workers that they should also activate the same environment of the master process with the `exeflags` option.
- We wrap the preamble into a `@everywhere begin ... end` block, and replace the for loop by a `pmap` call. We also add a `try ... catch` block to handle issues with specific files.

Here is the resulting script after the modifications:

```julia-auto
using Distributed

# add processes on the same machine
addprocs(4, topology=:master_worker, exeflags="--project=$(Base.active_project())")

# SETUP FOR ALL PROCESSES
# -----------------------
@everywhere begin
  # instantiate environment
  using Pkg; Pkg.instantiate()

  # load dependencies
  using CSV

  # HELPER FUNCTIONS
  # ----------------
  function process(infile, outfile)
    # read file from disk
    csv = CSV.read(infile)

    # perform calculations
    sleep(60) # pretend it takes time
    csv.new = rand(size(csv,1))

    # save new file to disk
    CSV.write(outfile, csv)
  end
end

# MAIN SCRIPT
# -----------

# relevant directories
indir = "data"
outdir = "results"

# files to process
infiles = readdir(indir, join=true)
outfiles = joinpath.(outdir, basename.(infiles))
nfiles = length(infiles)

status = pmap(1:nfiles) do i
  try
    process(infiles[i], outfiles[i])
    true # success
  catch e
    @warn "failed to process $(infiles[i])"
    false # failure
  end
end

```

We execute the script as before:

```shell
$ julia --project main.jl

```

### Questions

- Is there a more elegant method for the instantiation of the environment in remote workers? The `exeflags` feels like a hack.

* * *

## IO issues

Notice that we used `"data"` and `"results"` as our file paths in the script. If we try to run the script from outside the project directory (e.g. `proj`), we will get an error:

```shell
julia --project=proj proj/main.jl
ERROR: LoadError: SystemError: unable to read directory data: No such file or directory

```

Even worse, these file paths may not exist on different machines when we request multiple _remote_ workers. To solve this, we need to use paths relative to the `main.jl` source:

```julia-auto
# relevant directories
indir = joinpath(@ __DIR__ ,"data")
outdir = joinpath(@ __DIR__ ,"results")

```

or set these paths in the command line using some package like [DocOpt.jl](https://github.com/docopt/DocOpt.jl)

Our previous command should work with the suggested modifications:

```shell
julia --project=proj proj/main.jl

```

* * *

## Parallelization (remote machines)

Finally, we would like to run the script above in a cluster with hundreds of _remote_ worker processes. We don’t know in advance how many processes will be available because this is the job of a job scheduler (e.g. SLURM, PBS). We have the option of using [ClusterManagers.jl](https://github.com/JuliaParallel/ClusterManagers.jl) and the option to call the `julia` executable from a job script directly.

### Questions

- Could you please advise on the cleanest approach currently?
- How do you modify the script below to work with remote processes?

```julia-auto
using Distributed

# add processes on the same machine
addprocs(4, topology=:master_worker, exeflags="--project=$(Base.active_project())")

# SETUP FOR ALL PROCESSES
# -----------------------
@everywhere begin
  # instantiate environment
  using Pkg; Pkg.instantiate()

  # load dependencies
  using CSV

  # HELPER FUNCTIONS
  # ----------------
  function process(infile, outfile)
    # read file from disk
    csv = CSV.read(infile)

    # perform calculations
    sleep(60) # pretend it takes time
    csv.new = rand(size(csv,1))

    # save new file to disk
    CSV.write(outfile, csv)
  end
end

# MAIN SCRIPT
# -----------

# relevant directories
indir = joinpath(@ __DIR__ ,"data")
outdir = joinpath(@ __DIR__ ,"results")

# files to process
infiles = readdir(indir, join=true)
outfiles = joinpath.(outdir, basename.(infiles))
nfiles = length(infiles)

status = pmap(1:nfiles) do i
  try
    process(infiles[i], outfiles[i])
    true # success
  catch e
    @warn "failed to process $(infiles[i])"
    false # failure
  end
end

```

Appreciate if you can help improve this guide. I’ve created a repository on GitHub to track the improvements. Please feel free to submit PRs: [GitHub - Arpeggeo/julia-distributed-computing: The ultimate guide to distributed computing in Julia](https://github.com/juliohm/julia-distributed-computing)

### Contributors

@juliohm @samuel_okon

---

<div class="post-metadata">

**Author:** ![jishnub](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/jishnub/32/33620_2.png) [@jishnub](https://discourse.julialang.org/u/jishnub)\
**Post date:** [June 22, 2020, 3:43pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/2 "2020-06-22T15:43:41Z")

</div>

Perhaps an intro to `RemoteChannel`s for data transfer between workers is useful, as well as a guide to [ProgressMeter.jl](https://github.com/jekyllstein/ParallelProgressMeter.jl) for parallel jobs. Finally a guide to a parallel `mapreduce` operation using `RemoteChannel`s would be great.

---

<div class="post-metadata">

**Author:** ![dmolina](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/dmolina/32/5246_2.png) [@dmolina](https://discourse.julialang.org/u/dmolina)\
**Post date:** [June 22, 2020, 5:12pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/3 "2020-06-22T17:12:41Z")

</div>

Nice!

In respect to first Question, avoiding the execflags, you can avoiding it you put at the beginning

```julia
using Pkg

Pkg.activate(@ __DIR__ )

```

In that way, if the file is in a directory with a environment, it is always loaded, without considering the directory in which you are running Julia.

---

<div class="post-metadata">

**Author:** ![juliohm](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/juliohm/32/215266_2.png) [@juliohm](https://discourse.julialang.org/u/juliohm)\
**Post date:** [June 22, 2020, 5:30pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/4 "2020-06-22T17:30:44Z")

</div>

Thank you @jishnub for the feedback. I think these details can come later after we have the core of the distributed execution working out.

---

<div class="post-metadata">

**Author:** ![juliohm](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/juliohm/32/215266_2.png) [@juliohm](https://discourse.julialang.org/u/juliohm)\
**Post date:** [June 22, 2020, 5:31pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/5 "2020-06-22T17:31:32Z")

</div>

Thank you @dmolina, will try it locally. Nice suggestion.

---

<div class="post-metadata">

**Author:** ![juliohm](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/juliohm/32/215266_2.png) [@juliohm](https://discourse.julialang.org/u/juliohm)\
**Post date:** [June 22, 2020, 5:34pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/6 "2020-06-22T17:34:37Z")

</div>

It didn’t work for me @dmolina. We need to exeflags apparently.

---

<div class="post-metadata">

**Author:** ![dmolina](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/dmolina/32/5246_2.png) [@dmolina](https://discourse.julialang.org/u/dmolina)\
**Post date:** [June 22, 2020, 5:39pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/7 "2020-06-22T17:39:58Z")

</div>

It is trange, I did it for me, I had test it. Well, anyway, it is not too critical.

---

<div class="post-metadata">

**Author:** ![juliohm](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/juliohm/32/215266_2.png) [@juliohm](https://discourse.julialang.org/u/juliohm)\
**Post date:** [June 22, 2020, 5:45pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/8 "2020-06-22T17:45:22Z")

</div>

I’ve added the script to a repo so that other people can test it: [https://github.com/juliohm/julia-distributed-computing](https://github.com/juliohm/julia-distributed-computing)

Also updated the original post. @dmolina can you try your suggestion there and submit a PR if it works?

---

<div class="post-metadata">

**Author:** ![samuel\_okon](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/samuel_okon/32/10285_2.png) [@samuel\_okon](https://discourse.julialang.org/u/samuel_okon)\
**Post date:** [June 24, 2020, 1:24am UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/9 "2020-06-24T01:24:39Z")

</div>

@juliohm I feel the `@sync @everywhere ...` is redundant and should really be `@everywhere ...`.

---

<div class="post-metadata">

**Author:** ![juliohm](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/juliohm/32/215266_2.png) [@juliohm](https://discourse.julialang.org/u/juliohm)\
**Post date:** [June 24, 2020, 5:28pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/10 "2020-06-24T17:28:31Z")

</div>

Thank you @samuel_okon, could you please elaborate? You mean that the `@everywhere` is already blocking?

---

<div class="post-metadata">

**Author:** ![platawiec](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/platawiec/32/31914_2.png) [@platawiec](https://discourse.julialang.org/u/platawiec)\
**Post date:** [June 24, 2020, 5:51pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/11 "2020-06-24T17:51:22Z")

</div>

You can also connect to remote workers via a list of IPs, following is my typical startup which adds one worker per IP listed in each line of `machinefile.txt`:

```julia
machines = readlines("machinefile.txt")
machines_and_workers = [(machine, 1) for machine in machines]

addprocs(
    machines_and_workers,
    enable_threaded_blas=true,
    topology=:master_worker
)

```

I use an approach similar to the original post to instantiate the projects for each worker - needs to be said that the directory structure and contents must be identical between master and workers.

---

<div class="post-metadata">

**Author:** ![samuel\_okon](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/samuel_okon/32/10285_2.png) [@samuel\_okon](https://discourse.julialang.org/u/samuel_okon)\
**Post date:** [June 24, 2020, 7:08pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/12 "2020-06-24T19:08:14Z")

</div>

@juliohm yeah `@everywhere ...` calls `remotecall_eval` which is blocking

---

<div class="post-metadata">

**Author:** ![juliohm](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/juliohm/32/215266_2.png) [@juliohm](https://discourse.julialang.org/u/juliohm)\
**Post date:** [June 24, 2020, 7:23pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/13 "2020-06-24T19:23:07Z")

</div>

That is nice @platawiec, thank you for sharing. How do you generate the machine file? Are you relying on some job scheduler? Could you please share some more details on how we can automate the process?

---

<div class="post-metadata">

**Author:** ![juliohm](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/juliohm/32/215266_2.png) [@juliohm](https://discourse.julialang.org/u/juliohm)\
**Post date:** [June 24, 2020, 7:23pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/14 "2020-06-24T19:23:36Z")

</div>

Nice @samuel_okon, I’ve updated the instructions to remove the `@sync`.

---

<div class="post-metadata">

**Author:** ![juliohm](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/juliohm/32/215266_2.png) [@juliohm](https://discourse.julialang.org/u/juliohm)\
**Post date:** [June 24, 2020, 7:31pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/15 "2020-06-24T19:31:19Z")

</div>

@platawiec I wonder if we can use the `julia --machine-file` command line option and still enable the `exeflags` and `topology` options.

---

<div class="post-metadata">

**Author:** ![platawiec](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/platawiec/32/31914_2.png) [@platawiec](https://discourse.julialang.org/u/platawiec)\
**Post date:** [June 24, 2020, 8:01pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/16 "2020-06-24T20:01:49Z")

</div>

I’m not an expert, but I orchestrate and launch a cluster on AWS via KissCluster: [https://github.com/pszufe/KissCluster](https://github.com/pszufe/KissCluster). That takes care of everything after some setup. Workflow is basically to launch a bunch of workers, put their IPs in `machinefile`, make sure everything is configured correctly (permissions, directories, etc.), and start running.

There are other tools which allow you to emulate more traditional clusters in AWS, but if you don’t need the frills or you don’t like the overhead then KissCluster is a good lightweight alternative. Specifically, I wanted to bring up the `machinefile` option because my workflow doesn’t rely on `ClusterManagers.jl`.

---

<div class="post-metadata">

**Author:** ![juliohm](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/juliohm/32/215266_2.png) [@juliohm](https://discourse.julialang.org/u/juliohm)\
**Post date:** [June 24, 2020, 8:14pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/17 "2020-06-24T20:14:52Z")

</div>

Thank you @platawiec. I think I know how to make it work with ClusterManagers.jl on general clusters. I will give it a try soon, and will report the results.

---

<div class="post-metadata">

**Author:** ![healyp](https://avatars.discourse-cdn.com/v4/letter/h/67e7ee/32.png) [@healyp](https://discourse.julialang.org/u/healyp)\
**Post date:** [June 24, 2020, 9:12pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/18 "2020-06-24T21:12:55Z")

</div>

The following hacky function can take the place of a `machinefile`. Perhaps a little too much of a corner case but I often want to fire up a set of jobs on as many of our dual-boot student lab machines as I can.

```julia
function whatsup()
    cs144l = map(i->@sprintf("cs144l-%02d.csis.ul.ie", i), 1:40)
    cs244l = map(i->@sprintf("cs244l-%02d.csis.ul.ie", i), 1:40)
    cs305l = map(i->@sprintf("cs305l-%02d.csis.ul.ie", i), [x for x in 1:32 if x != 7])# cs305l-07 misconfigured

    possibles = vcat(cs305l, cs144l, cs244l)
    usables = []

    # https://discourse.julialang.org/t/how-can-i-ping-an-ip-adress-in-julia/3380
    # with v1.0 modifications
    @sync for p in possibles
        @async push!(usables, match(r"PING ([^\s]*)\s",
                                    split(read(`ping -c 1 $p`, String), "\n";
                                          keepempty=false)[1])[1]
                     )
        sleep(0.1)
    end;
    string.(usables)
end

```

Edited as per suggestion of @tkf.

---

<div class="post-metadata">

**Author:** ![juliohm](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/juliohm/32/215266_2.png) [@juliohm](https://discourse.julialang.org/u/juliohm)\
**Post date:** [June 24, 2020, 9:55pm UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/19 "2020-06-24T21:55:00Z")

</div>

Thank you @healyp for sharing. Handy funcion when you know beforehand the IPs.

---

<div class="post-metadata">

**Author:** ![tkf](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/tkf/32/17635_2.png) [@tkf](https://discourse.julialang.org/u/tkf)\
**Post date:** [June 25, 2020, 6:15am UTC](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867/20 "2020-06-25T06:15:27Z")

</div>

`@spawn push!(usables, ...)` has a data race. Use `@async` instead. See also [the manual](https://docs.julialang.org/en/v1.6-dev/manual/multi-threading/#Data-race-freedom-1):

> You are entirely responsible for ensuring that your program is data-race free, and nothing promised here can be assumed if you do not observe that requirement.

You also need `@sync for p in possibles` to wait for all the tasks. With the code as-is, it’s very possible that the function returns before any of the task finishes.

[Next page](https://discourse.julialang.org/t/the-ultimate-guide-to-distributed-computing/41867.md?page=2)
