Julia muitas alocações usando Distributed and SharedArrays com @ sync / @ async

Sep 11 2020

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

2 PrzemyslawSzufel Sep 11 2020 at 05:57

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 executar addprocsem alguma parte do script que realiza a inicialização ou talvez os processos sejam adicionados iniciando o juliaprocesso 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