# Julia vs Python's Dask: Known speed comparisons?

**URL:** <https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459>\
**Category:** Julia at Scale\
**Tags:** question, parallel, distributed\
**Created:** [October 3, 2019, 6:03pm UTC](https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459 "2019-10-03T18:03:47Z")\
**Posts on this page:** 14\
**Page:** 1

<div class="post-metadata">

**Author:** ![joshhjacobson](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/joshhjacobson/32/10578_2.png) [@joshhjacobson](https://discourse.julialang.org/u/joshhjacobson)\
**Post date:** [October 3, 2019, 6:03pm UTC](https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459/1 "2019-10-03T18:03:47Z")

</div>

Hi all,

I’m new to the Julia community and have arrived in hopes of improving the speed of parallelization tasks I’ve been running with Dask (distributed) in Python. I know there are some Julia packages that are inspired by Dask (e.g., [Dagger.jl](https://github.com/JuliaParallel/Dagger.jl)), but my main question is whether anyone is aware of direct comparisons between the speed of Dask and similar implementations in Julia?

Thanks,  
Josh

---

<div class="post-metadata">

**Author:** ![xiaodai](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/xiaodai/32/15937_2.png) [@xiaodai](https://discourse.julialang.org/u/xiaodai)\
**Post date:** [October 3, 2019, 10:56pm UTC](https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459/2 "2019-10-03T22:56:16Z")

</div>

[https://diskframe.com/articles/vs-dask-juliadb.html](https://diskframe.com/articles/vs-dask-juliadb.html)

---

<div class="post-metadata">

**Author:** ![zgornel](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/zgornel/32/217487_2.png) [@zgornel](https://discourse.julialang.org/u/zgornel)\
**Post date:** [October 4, 2019, 1:51pm UTC](https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459/3 "2019-10-04T13:51:52Z")

</div>

If you mean the speed of parsing a DAG and generating tasks graphs, that is mostly negligible (compared to the time the tasks take themselves). You can check [Dispatcher.jl](https://github.com/invenia/Dispatcher.jl) as well.

---

<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:** [October 4, 2019, 2:27pm UTC](https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459/4 "2019-10-04T14:27:27Z")

</div>

Dagger works, but its performance is still not optimal in certain common cases. If your operations can safely be run in true parallel with multithreading, then you can use the Dagger master branch and instruct the scheduler to multithread execution of DAG tasks on each worker process. I haven’t benchmarked the performance gains from this yet, but if you have an MWE that’s close to what you’re working on, I could probably use it to tune Dagger’s scheduler to try to beat Dask (assuming Dagger isn’t already faster, which it probably isn’t right now).

---

<div class="post-metadata">

**Author:** ![joshhjacobson](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/joshhjacobson/32/10578_2.png) [@joshhjacobson](https://discourse.julialang.org/u/joshhjacobson)\
**Post date:** [October 4, 2019, 4:11pm UTC](https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459/5 "2019-10-04T16:11:31Z")

</div>

Here’s a general idea of what I’m doing (though not a working example). In my specific case, I’m working with a 12x2000x2000 array (~70mb), and constructing a 750x2000x2000 array by looping over the second and third dimensions (embarrassingly parallel in that results at a given index are not dependent on any other index). Using 32 cores, the Dask implementation looks like:

```python

import numpy as np
import dask
from dask.distributed import LocalCluster, Client

def loop_func(my_array):
        # Get the shape
        nb, ni, nj = my_array.shape

        new_array = np.full((750, ni, nj), np.nan, dtype=np.float32)

        for i in range(ni):
            for j in range(nj):
                if (not (np.isnan(my_array[:, i, j]).any())):
                    mean_params = my_array[:4, i, j]
                    cov_mat = my_array[4:, i, j].reshape((3,3))
                    new_array[:, i, j] = do_stuff(mean_params, cov_mat)

        return new_array
        
## dask
# setup dask cluster 

def call_dask():
    parm_sims = dask.array.map_blocks(sim_parms_uq, param_array, num_pars, num_uq_its, scale_parm_index,
                                      chunks=(-1, chunk_size, chunk_size), dtype=np.float32)
    parm_sims.compute()
    return
    
print(timeit.Timer(call_dask).repeat(3, 5))

# shutdown dask cluster

```

I’ve also come up with the following set of modules in Julia which are called through Python using [PyJulia](https://github.com/JuliaPy/pyjulia) and process the same array that is passed to Dask:

```julia

module distributed_utils
export loop_func!

using DistributedArrays

function loop_func!(new_array, A)
    A_loc = localpart(A)
    sim_dat = fill(NaN, size(localpart(new_array)))
    
    for j = 1:size(A_loc, 3), i = 1:size(A_loc, 2)
        if !any(isnan, A_loc[:, i, j])
            mean_params = A_loc[1:3, i, j]
            cov_mat = transpose(reshape(A_loc[4:end, i, j], 3, 3))
            cov_mat = convert(Array, cur_cov_mat)
            if isposdef(cov_mat)
                sim_dat[:, i, j] = do_stuff(mean_params, cov_mat)
            end
        end
    end
    new_array[:L] = sim_dat
end;

end; #module

module my_module
export python_wrapper

using Distributed
CPU_CORES = length(Sys.cpu_info())
addprocs(CPU_CORES - 1)

module_dir = "/home/jovyan/workdir/lib"
@everywhere push!(LOAD_PATH, $module_dir)
@everywhere using distributed_utils: loop_func!
@everywhere using DistributedArrays

function f_distributed!(new_array, my_array)
    @sync begin
        for p in procs(new_array)
            @async remotecall_wait(loop_func!, p, new_array, my_array)
        end
    end
end;

function python_wrapper(dataset)
    ## NOTE: these arrays are "chunked" automatically (could be done manually though)
    my_array = distribute(dataset)
    new_array = dfill(NaN, 750, size(dataset, 2), size(dataset, 3))

    f_distributed!(new_array, my_array)
    
    new_array = convert(Array, new_array)
    return new_array
end;

end; # module

```

Then in Python

```python

import julia
jl = julia.Julia(compiled_modules=False)
jl.include('./lib/my_module.jl')
python_wrapper = julia.Main.my_module.python_wrapper

def call_julia():
    python_wrapper(my_array)

## call once to compile
call_julia()
    
print(timeit.Timer(call_julia).repeat(3, 5))

```

Interestingly, the overall process is a bit slower in Julia. The results from Python’s `timeit` module are

- Dask: 26 seconds (average)
- Julia: 32 seconds (average)

I was surprised by these results and have to wonder if I should be doing something different in the Julia implementation, or if Dask might really be faster in this case? One thought I have is that the Julia implementation has a check for positive definiteness in the double for-loop and the Python version does not. However, the matrices being checked are small (3x3) so I’m not sure this would be significant.

---

<div class="post-metadata">

**Author:** ![ihnorton](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/ihnorton/32/26_2.png) [@ihnorton](https://discourse.julialang.org/u/ihnorton)\
**Post date:** [October 4, 2019, 6:00pm UTC](https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459/6 "2019-10-04T18:00:34Z")

</div>

Out of curiosity, have you benchmarked this problem with single-threaded Julia, or tried [multi-threading](https://julialang.org/blog/2019/07/multithreading) rather than using DistributedArrays? Since you are on a single 32 core machine (per `import LocalCluster`), multi-threading might be faster.

---

<div class="post-metadata">

**Author:** ![Ratingulate](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/ratingulate/32/9242_2.png) [@Ratingulate](https://discourse.julialang.org/u/Ratingulate)\
**Post date:** [October 4, 2019, 7:02pm UTC](https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459/7 "2019-10-04T19:02:00Z")

</div>

Dask’s distributed scheduler is highly optimized and stress tested, so it could just be that Julia’s needs a bit more work.

---

<div class="post-metadata">

**Author:** ![Elrod](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/elrod/32/22461_2.png) [@Elrod](https://discourse.julialang.org/u/Elrod)\
**Post date:** [October 5, 2019, 2:58am UTC](https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459/8 "2019-10-05T02:58:09Z")

</div>

Using views or, even better, `StaticArray`s would likely speed up the code.  
If going the `StaticArray`s route, be sure not to allocate any `Array` objects through slicing, etc.  
I’d also recommend double checking Julia’s performance tips. Eg, checking for type instabilities. A heuristic I often use is that if a serial version of the code is low on allocations, it’s probably pretty okay and not doing anything unintentionally silly.

---

<div class="post-metadata">

**Author:** ![zgornel](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/zgornel/32/217487_2.png) [@zgornel](https://discourse.julialang.org/u/zgornel)\
**Post date:** [October 5, 2019, 11:37am UTC](https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459/9 "2019-10-05T11:37:44Z")

</div>

It would be nice if Rocklin would move to the julia side 😉

---

<div class="post-metadata">

**Author:** ![Ratingulate](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/ratingulate/32/9242_2.png) [@Ratingulate](https://discourse.julialang.org/u/Ratingulate)\
**Post date:** [October 6, 2019, 5:37pm UTC](https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459/10 "2019-10-06T17:37:14Z")

</div>

Indeed it would be

---

<div class="post-metadata">

**Author:** ![joshhjacobson](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/joshhjacobson/32/10578_2.png) [@joshhjacobson](https://discourse.julialang.org/u/joshhjacobson)\
**Post date:** [October 7, 2019, 4:42pm UTC](https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459/11 "2019-10-07T16:42:22Z")

</div>

> [@ihnorton](#):
>
> Out of curiosity, have you benchmarked this problem with single-threaded Julia, or tried [multi-threading](https://julialang.org/blog/2019/07/multithreading) rather than using DistributedArrays? Since you are on a single 32 core machine (per `import LocalCluster`), multi-threading might be faster.

I have not tried multi-threading, but will give it a shot. Thanks for the suggestion!

While I am on a single machine for now, this is mainly a proof of concept. If I were to scale this to multiple machines, would you still suggest multi-threading?

---

<div class="post-metadata">

**Author:** ![joshhjacobson](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/joshhjacobson/32/10578_2.png) [@joshhjacobson](https://discourse.julialang.org/u/joshhjacobson)\
**Post date:** [October 7, 2019, 4:46pm UTC](https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459/12 "2019-10-07T16:46:59Z")

</div>

> Using views or, even better, `StaticArray` s would likely speed up the code.

Thanks for pointing out the `StaticArray`s package, I hadn’t come across it until now. I’ll check it out!

---

<div class="post-metadata">

**Author:** ![ihnorton](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/ihnorton/32/26_2.png) [@ihnorton](https://discourse.julialang.org/u/ihnorton)\
**Post date:** [October 11, 2019, 1:46am UTC](https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459/13 "2019-10-11T01:46:23Z")

</div>

> [@joshhjacobson](#):
>
> While I am on a single machine for now, this is mainly a proof of concept. If I were to scale this to multiple machines, would you still suggest multi-threading?

It might be useful on a per-node basis, although I don’t know how it composes with operations on DistributedArrays.

The reason I asked about single-threaded benchmarks is that I am curious if you have a baseline for the performance of the python vs julia implementations.

---

<div class="post-metadata">

**Author:** ![joshhjacobson](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/joshhjacobson/32/10578_2.png) [@joshhjacobson](https://discourse.julialang.org/u/joshhjacobson)\
**Post date:** [October 15, 2019, 3:58pm UTC](https://discourse.julialang.org/t/julia-vs-pythons-dask-known-speed-comparisons/29459/14 "2019-10-15T15:58:16Z")

</div>

I’m not sure I understand your question, do you mean benchmarks for serial implementations? For the results above, I was using a single thread per worker in the Dask implementation, and as far as I know Julia would’ve done the same.
