Julia: paralléliser les opérations sur des structures de données complexes (ex: DataFrames)

Sep 08 2020

Je souhaite traiter un certain nombre de grands ensembles de données en parallèle. Malheureusement, l'accélération que j'obtiens en utilisant Threads.@threadsest très sous-linéaire, comme le montre l'exemple simplifié suivant.

(Je suis très nouveau avec Julia, alors excuses si j'ai raté quelque chose d'évident)

Créons des données d'entrée factices - 8 dataframes avec 2 colonnes entières chacune et 10 millions de lignes:

using DataFrames

n = 8
dfs = Vector{DataFrame}(undef, n)
for i = 1:n
    dfs[i] = DataFrame(Dict("x1" => rand(1:Int64(1e7), Int64(1e7)), "x2" => rand(1:Int64(1e7), Int64(1e7))))
end

Maintenant, faites un peu de traitement sur chaque dataframe (group by x1and sum x2)

function process(df::DataFrame)::DataFrame
    combine([:x2] => sum, groupby(df, :x1))
end

Enfin, comparez la vitesse de traitement sur une seule trame de données avec celle sur les 8 trames de données en parallèle. La machine sur laquelle je l'exécute a 50 cœurs et Julia a été démarrée avec 50 threads, donc idéalement, il ne devrait pas y avoir beaucoup de décalage horaire.

julia> dfs_res = Vector{DataFrame}(undef, n)

julia> @time for i = 1:1
           dfs_res[i] = process(dfs[i])
       end
  3.041048 seconds (57.24 M allocations: 1.979 GiB, 4.20% gc time)

julia> Threads.nthreads()
50

julia> @time Threads.@threads for i = 1:n
           dfs_res[i] = process(dfs[i])
       end
  5.603539 seconds (455.14 M allocations: 15.700 GiB, 39.11% gc time)

Ainsi, l'exécution en parallèle prend presque deux fois plus de temps par ensemble de données (cela s'aggrave avec plus d'ensembles de données). J'ai le sentiment que cela a quelque chose à voir avec une gestion de la mémoire inefficace. Le temps du GC est assez élevé pour la deuxième manche. Et je suppose que la préallocation avec undefn'est pas efficace pour l' DataFrameart. Presque tous les exemples que j'ai vus pour le traitement parallèle dans Julia sont réalisés sur des tableaux numériques avec des tailles fixes et connues a priori. Cependant, ici, les ensembles de données peuvent avoir des tailles, des colonnes arbitraires, etc. Dans R, des flux de travail comme celui-ci peuvent être réalisés très efficacement avec mclapply. Y a-t-il quelque chose de similaire (ou un modèle différent mais efficace) chez Julia? J'ai choisi d'utiliser les threads et non le multi-traitement pour éviter de copier des données (Julia ne semble pas prendre en charge le modèle de processus de fourche comme R / mclapply).

Réponses

1 PrzemyslawSzufel Sep 08 2020 at 18:03

Le multithreading dans Julia ne s'adapte pas bien au-delà des 16threads. Par conséquent, vous devez utiliser le multitraitement à la place. Votre code pourrait ressembler à ceci:

using DataFrames, Distributed
addprocs(4) # or 50
@everywhere using DataFrames, Distributed

n = 8
dfs = Vector{DataFrame}(undef, n)
for i = 1:n
    dfs[i] = DataFrame(Dict("x1" => rand(1:Int64(1e7), Int64(1e7)), "x2" => rand(1:Int64(1e7), Int64(1e7))))
end

@everywhere function process(df::DataFrame)::DataFrame
    combine([:x2] => sum, groupby(df, :x1))
end

dfs_res = @distributed (vcat) for i = 1:n
      df = process(dfs[i])
      (i, myid(), df)
end

Ce qui est important dans ce type de code, c'est que le transfert de données entre les processus prend du temps. Alors parfois, vous voudrez peut-être simplement garder des DataFrames séparés sur des travailleurs distincts. Comme toujours - cela dépend de votre architecture de traitement.

Modifier quelques notes sur la performance

Pour les tests, ayez votre code dans les fonctions et utilisez consts (ou utilisez BenchamrTools.jl)

using DataFrames

const dfs = [DataFrame(Dict("x1" => rand(1:Int64(1e7), Int64(1e7)), "x2" => rand(1:Int64(1e7), Int64(1e7)))) for i in 1:8 ]

function process(df::DataFrame)::DataFrame
    combine([:x2] => sum, groupby(df, :x1))
end

function p1!(res, d)
    for i = 1:8
        res[i] = process(dfs[i])
    end
end


function p2!(res, d)
     Threads.@threads for i = 1:8
        res[i] = process(dfs[i])
    end
end

const dres = Vector{DataFrame}(undef, 8)

Et voici le résultat

julia> GC.gc();@time p1!(dres, dfs)
 30.840718 seconds (507.28 M allocations: 16.532 GiB, 6.42% gc time)

julia> GC.gc();@time p1!(dres, dfs)
 30.827676 seconds (505.66 M allocations: 16.451 GiB, 7.91% gc time)

julia> GC.gc();@time p2!(dres, dfs)
 18.002533 seconds (505.77 M allocations: 16.457 GiB, 23.69% gc time)

julia> GC.gc();@time p2!(dres, dfs)
 17.675169 seconds (505.66 M allocations: 16.451 GiB, 23.64% gc time)

Pourquoi la différence n'est qu'environ 2x sur une machine à 8 cœurs - parce que nous avons passé la plupart du temps à ramasser des ordures! (regardez la sortie dans votre question - le problème est le même) Lorsque vous utilisez moins de RAM, vous verrez une meilleure vitesse de multithreading jusqu'à 3x.