Saltar a contenido

Para paralelizar tareas vamos a usar el enfoque de que vamos a dividir el problema en tareas paralelas hasta que se llega a cierto tamaño o condicion, para el caso del mergeSort se va hablar de una profundidad

// Sesion 12 - Paralelismo de tareas  
// MergeSort paralelo. Se ordenan en paralelo las dos mitades  
// con parallel(a,b); la mezcla queda secuencial.  
// Recursion estructural sobre listas (head/tail/isEmpty);  
// la mezcla auxiliar es de cola con un acumulador invertido.  
package taller  

import common._  

class MergeSortPar {  
  def merge(xs: List[Int], ys: List[Int]): List[Int] = {  
    def aux(as: List[Int], bs: List[Int], acc: List[Int]): List[Int] =  
      if (as.isEmpty) acc.reverse ::: bs  
      else if (bs.isEmpty) acc.reverse ::: as  
      else if (as.head <= bs.head) aux(as.tail, bs, as.head :: acc)  
      else aux(as, bs.tail, bs.head :: acc)  

    aux(xs, ys, Nil)  
  }  

  def partir(xs: List[Int]): (List[Int], List[Int]) = {  
    val n = xs.length  
    val m = n / 2  

    def aux(k: Int, ys: List[Int], izq: List[Int]): (List[Int], List[Int]) =  
      if (k <= 0) (izq.reverse, ys)  
      else aux(k - 1, ys.tail, ys.head :: izq)  

    aux(m, xs, Nil)  
  }  

  def mergeSortSec(xs: List[Int]): List[Int] =  
    if (xs.isEmpty || xs.tail.isEmpty) xs  
    else {  
      val (izq, der) = partir(xs)  
      merge(mergeSortSec(izq), mergeSortSec(der))  
    }  

  def mergeSortPar(profMax: Int)(xs: List[Int])(prof: Int = 0): List[Int] =  
    if (xs.isEmpty || xs.tail.isEmpty) xs  
    else if (prof >= profMax) mergeSortSec(xs)  
    else {  
      val (izq, der) = partir(xs)  
      val (izqOrd, derOrd) = parallel(  
        mergeSortPar(profMax)(izq)(prof + 1),  
        mergeSortPar(profMax)(der)(prof + 1)  
      )  
      merge(izqOrd, derOrd)  
    }  

  /*@main def demoMergeSortPar(): Unit = {  
    val rng = new scala.util.Random(42)    val xs = List.fill(10_000)(rng.nextInt(100_000))    val sec = mergeSortSec(xs)    val par = mergeSortPar(profMax = 4)(0)(xs)    assert(sec == par, "discrepancia entre secuencial y paralelo")    assert(sec == xs.sorted, "no quedaron ordenados")    println(s"mergeSort: ${xs.size} elementos ordenados ok")  }*/}
/*  
 * This Scala source file was generated by the Gradle 'init' task. */package taller  


import org.scalameter._  
import common._  


object App {  
  def main(args: Array[String]): Unit = {  
    val objMergeSortPar = new MergeSortPar()  
    val randr = scala.util.Random  
    randr.setSeed(42)  
    val arr:List[Int] = (1 to 100000).toList.map(x => randr.nextInt())  

    val tseq = withWarmer(new Warmer.Default) measure{  
      objMergeSortPar.mergeSortSec(arr)  
    }  
    val tpar2 = withWarmer(new Warmer.Default) measure{  
      objMergeSortPar.mergeSortPar(2)(arr)(0)  
    }  
    val tpar4 = withWarmer(new Warmer.Default) measure{  
      objMergeSortPar.mergeSortPar(4)(arr)(0)  
    }  
    val tpar5 = withWarmer(new Warmer.Default) measure{  
      objMergeSortPar.mergeSortPar(5)(arr)(0)  
    }  

    println("Secuencial "+tseq)  
    println("Paralelo prof 2 "+tpar2)  
    println("Paralelo prof 4 "+tpar4)  
    println("Paralelo prof 5 "+tpar5)  
  }  
  def greeting(): String = "Hello, world!"  
}

Al ejecutar encontramos lo siguiente

Secuencial 52.917458 ms
Paralelo prof 2 28.813115 ms
Paralelo prof 4 26.907541 ms
Paralelo prof 5 31.255206 ms

Esto quiere decir que se obtiene una mejora hasta profundidad 4, despues de eso tenemos overhead (por gestion de hilos)

Scan

Es la operación que retorna los valores del acumulador cuando hago reduce

scala> List(1,2,3).foldLeft(0)((acc,x) => acc + x)
val res7: Int = 6

scala> List(1,2,3).scanLeft
(0)((acc,x) => acc + x)
val res8: List[Int] = List(0, 1, 3, 6)

scala> List(1,2,3).scanRight(0)((acc,x) => acc + x)
val res9: List[Int] = List(6, 5, 3, 0)

La implemetnacion paralela es

// Sesion 12 - Paralelismo de tareas
// scanLeft paralelo. Idea: combinar map + reduce reutilizando
// resultados intermedios. Esta version trabaja sobre arboles
// de expresion: en una primera pasada (upsweep) se calcula el
// reduce de cada subarbol; en una segunda pasada (downsweep)
// se distribuye el prefijo desde la raiz.

// Arbol con anotacion de reduce parcial en cada nodo.
sealed trait ArbolRed[+A]
case class HojaRed[+A](valor: A, reducido: A) extends ArbolRed[A]
case class NodoRed[+A](izq: ArbolRed[A], der: ArbolRed[A], reducido: A) extends ArbolRed[A]

def reducidoDe[A](t: ArbolRed[A]): A = t match {
  case HojaRed(_, r) => r
  case NodoRed(_, _, r) => r
}

// Upsweep: anotar cada subarbol con el reduce de sus hojas.
def upsweep[A](t: Arbol[A])(f: (A, A) => A): ArbolRed[A] = t match {
  case Hoja(v) => HojaRed(v, v)
  case Nodo(l, r) =>
    val (la, ra) = parallel(upsweep(l)(f), upsweep(r)(f))
    NodoRed(la, ra, f(reducidoDe(la), reducidoDe(ra)))
}

// Downsweep: bajar el prefijo a izquierda y derecha.
// Pre: a0 es el prefijo (acumulador) que llega desde la raiz.
def downsweep[A](t: ArbolRed[A], a0: A, f: (A, A) => A): Arbol[A] = t match {
  case HojaRed(v, _) => Hoja(f(a0, v))
  case NodoRed(l, r, _) =>
    val (lp, rp) = parallel(
      downsweep(l, a0, f),
      downsweep(r, f(a0, reducidoDe(l)), f)
    )
    Nodo(lp, rp)
}

def scanPar[A](t: Arbol[A], a0: A)(f: (A, A) => A): Arbol[A] =
  downsweep(upsweep(t)(f), a0, f)

@main def demoScanPar(): Unit = {
  // Hojas: 1, 2, 3, 4. Con f = +, scan con a0=0 produce
  // prefijos parciales sumando lo que viene a la izquierda.
  val t: Arbol[Int] = Nodo(Nodo(Hoja(1), Hoja(2)), Nodo(Hoja(3), Hoja(4)))
  val s = scanPar(t, 0)((a, b) => a + b)
  // hoja(1) -> 0+1 = 1
  // hoja(2) -> 0+1+2 = 3
  // hoja(3) -> 0+(1+2)+3 = 6
  // hoja(4) -> 0+(1+2+3)+4 = 10
  val esperado: Arbol[Int] = Nodo(Nodo(Hoja(1), Hoja(3)), Nodo(Hoja(6), Hoja(10)))
  assert(s == esperado, s"scanPar fallo: $s")
  println(s"scanPar produjo $s")
}