Julia muitas alocações usando Distributed and SharedArrays com @ sync / @ async
Estou tentando entender como usar o pacote Distribuído junto com SharedArrays para realizar operações paralelas com julia. Apenas como exemplo, estou usando um método simples de média Montecarlo
using Distributed
using SharedArrays
using Statistics
const NWorkers = 2
const Ns = Int(1e6)
function parallelRun()
addprocs(NWorkers)
procsID = workers()
A = SharedArray{Float64,1}(Ns)
println("starting loop")
for i=1:2:Ns
#parallel block
@sync for p=1:NWorkers
@async A[i+p-1] = remotecall_fetch(rand,procsID[p]);
end
end
println(mean(A))
end
function singleRun()
A = zeros(Ns)
for i=1:Ns
A[i] = rand()
end
println(mean(A))
end
No entanto, se eu @time ambas as funções eu obtenho
julia> @time singleRun()
0.49965531193003165
0.009762 seconds (17 allocations: 7.630 MiB)
julia> @time parallelRun()
0.4994892300029917
46.319737 seconds (66.99 M allocations: 2.665 GiB, 1.01% gc time)
Em particular, há muito mais alocações na versão paralela, o que torna o código muito mais lento.
Estou esquecendo de algo?
A propósito, o motivo pelo qual estou usando @sync e @async (mesmo se não for necessário neste quadro, uma vez que cada amostra pode ser calculada em ordem aleatória) é apenas porque eu gostaria de aplicar a mesma estratégia para resolver um PDE parabólico numericamente com algo na linha de
for t=1:time_steps
#parallel block
@sync for p=1:NWorkers
@async remotecall(make_step_PDE,procsID[p],p);
end
end
onde cada trabalhador indexado por p deve trabalhar em um conjunto separado de índices da minha equação.
desde já, obrigado
Respostas
Existem os seguintes problemas no seu código:
- Você está gerando uma tarefa remota para cada valor de
ia e isso é caro e, no final, leva muito tempo. Basicamente, a regra é usar@distributedmacro para o balanceamento de carga entre os workers, o que irá compartilhar o trabalho uniformemente. - Nunca coloque
addprocsdentro de sua função de trabalho porque toda vez que você a executa, toda vez que adiciona novos processos - gerar um novo processo Julia também leva muito tempo e foi incluído em suas medições. Na prática, isso significa que você deseja executaraddprocsem alguma parte do script que realiza a inicialização ou talvez os processos sejam adicionados iniciando ojuliaprocesso com-pou--machine-fileparâmetro - Por fim, execute
@timesempre sempre duas vezes - na primeira medição@timetambém está medindo os tempos de compilação e a compilação em um ambiente distribuído leva muito mais tempo do que em um único processo.
Sua função deve ser mais ou menos assim
using Distributed, SharedArrays
addprocs(4)
@everywhere using Distributed, SharedArrays
function parallelRun(Ns)
A = SharedArray{Float64,1}(Ns)
@sync @distributed for i=1:Ns
A[i] = rand();
end
println(mean(A))
end
Você também pode considerar a divisão completa dos dados entre os trabalhadores. Em alguns cenários, isso é menos sujeito a bugs e permite que você distribua por muitos nós:
using Distributed, DistributedArrays
addprocs(4)
@everywhere using Distributed, DistributedArrays
function parallelRun2(Ns)
d = dzeros(Ns) #creates an array distributed evenly on all workers
@sync @distributed for i in 1:Ns
p = localpart(d)
p[((i-1) % Int(Ns/nworkers())+1] = rand()
end
println(mean(d))
end