Julia: parallelizza le operazioni su strutture dati complesse (es. DataFrame)

Sep 08 2020

Vorrei elaborare in parallelo una serie di grandi set di dati. Sfortunatamente la velocità che sto ottenendo dall'utilizzo Threads.@threadsè molto sublineare, come mostra il seguente esempio semplificato.

(Sono molto nuovo su Julia, quindi mi scuso se mi sono perso qualcosa di ovvio)

Creiamo alcuni dati di input fittizi: 8 dataframe con 2 colonne intere ciascuna e 10 milioni di righe:

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

Ora esegui alcune elaborazioni su ogni dataframe (raggruppa per x1e somma x2)

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

Infine, confronta la velocità di elaborazione su un singolo dataframe con l'esecuzione su tutti gli 8 dataframe in parallelo. La macchina su cui sto eseguendo questo ha 50 core e Julia è stata avviata con 50 thread, quindi idealmente non dovrebbe esserci molta differenza di tempo.

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)

Quindi l'esecuzione parallela richiede quasi il doppio del tempo per set di dati (questo peggiora con più set di dati). Ho la sensazione che questo abbia qualcosa a che fare con una gestione inefficiente della memoria. Il tempo GC è piuttosto alto per la seconda manche. E presumo che la preallocazione con undefnon sia efficiente per DataFrames. Quasi tutti gli esempi che ho visto per l'elaborazione parallela in Julia sono eseguiti su array numerici con dimensioni fisse e note a priori. Tuttavia qui i set di dati potrebbero avere dimensioni arbitrarie, colonne, ecc. In flussi di lavoro R come questo può essere fatto in modo molto efficiente con mclapply. C'è qualcosa di simile (o uno schema diverso ma efficiente) in Julia? Ho scelto di utilizzare thread e non multielaborazione per evitare di copiare i dati (Julia non sembra supportare il modello di processo fork come R / mclapply).

Risposte

1 PrzemyslawSzufel Sep 08 2020 at 18:03

Il multithreading in Julia non scala ben oltre i 16thread. Quindi è necessario utilizzare invece il multiprocessing. Il tuo codice potrebbe essere simile a questo:

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

Ciò che è importante in questo tipo di codice è che il trasferimento dei dati tra i processi richiede tempo. Quindi a volte potresti voler mantenere messaggi separati DataFramesu lavoratori separati. Come sempre, dipende dalla tua architettura di elaborazione.

Modifica alcune note sulla performance

Per i test, constinserisci il tuo codice in functions e usa s (o usa 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)

E qui risultato

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)

Perché la differenza è solo circa 2x su una macchina a 8 core - perché abbiamo passato la maggior parte del tempo a raccogliere i rifiuti! (guarda l'output nella tua domanda: il problema è lo stesso) Quando usi meno RAM vedrai una migliore velocità di multithreading fino a 3x.