Julia: parallelizza le operazioni su strutture dati complesse (es. DataFrame)
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
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.