Julia muchas asignaciones usando Distributed y SharedArrays con @ sync / @ async

Sep 11 2020

Estoy tratando de entender cómo usar el paquete Distribuido junto con SharedArrays para realizar operaciones paralelas con julia. Solo como ejemplo, estoy tomando un método simple de promedio de 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

Sin embargo, si @time ambas funciones obtengo

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)

En particular, hay muchas más asignaciones en la versión paralela, lo que hace que el código sea mucho más lento.

¿Me estoy perdiendo de algo?

Por cierto, la razón por la que estoy usando @sync y @async (incluso si no es necesario en este marco, ya que cada muestra se puede calcular en orden aleatorio) es solo porque me gustaría aplicar la misma estrategia para resolver un PDE parabólico numéricamente con algo en la línea de

    for t=1:time_steps

        #parallel block
        @sync for p=1:NWorkers
                @async remotecall(make_step_PDE,procsID[p],p);
        end
    end

donde cada trabajador indexado por p debería trabajar en un conjunto disjunto de índices de mi ecuación.

Gracias por adelantado

Respuestas

2 PrzemyslawSzufel Sep 11 2020 at 05:57

Existen los siguientes problemas en su código:

  • Está generando una tarea remota para cada valor de ia y esto es caro y, al final, lleva mucho tiempo. Básicamente, la regla general es usar @distributedmacro para el equilibrio de carga entre los trabajadores, esto solo compartirá el trabajo de manera uniforme.
  • Nunca ponga addprocsdentro su función de trabajo porque cada vez que la ejecuta, cada vez que agrega nuevos procesos, generar un nuevo proceso de Julia también lleva mucho tiempo y esto se incluyó en sus mediciones. En la práctica, esto significa que desea ejecutar addprocsen alguna parte del script que realiza la inicialización o quizás los procesos se agregan iniciando el juliaproceso con -po --machine-fileparámetro
  • Por último, ejecute @timesiempre siempre dos veces: en la primera medición @timetambién se miden los tiempos de compilación y la compilación en un entorno distribuido lleva mucho más tiempo que en un solo proceso.

Tu función debería verse más o menos así

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

También podría considerar dividir completamente los datos entre los trabajadores. Esto en algunos escenarios es menos propenso a errores y le permite distribuir en muchos nodos:

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