# Creamos el Spark Context

In [143]:
from pyspark import SparkContext
from pyspark.sql import SQLContext
from pydrive.drive import GoogleDrive
#from google.colab import auth
from pydrive.auth import GoogleAuth
from oauth2client.client import GoogleCredentials

#auth.authenticate_user()
#gauth = GoogleAuth()
#gauth.credentials = GoogleCredentials.get_application_default()
#drive = GoogleDrive(gauth)


#  Por si hay mas de un contexto de PySpark corriendo (por ejemplo, otro Notebook), esto para utilizar el mismo.
sc = SparkContext.getOrCreate()

In [144]:
type(sc)

pyspark.context.SparkContext

# Lectura de datos en Spark

## Paralelizando una coleccion de python

In [145]:
## creamos 1000 enteros en una lista
integersList = range(1,1001)
len(integersList)

1000

In [146]:
## Paralelizamos la coleccion utilizando 8 particiones o slices
## Esta operacion es una transformacion de datos en un RDD
## Dado que Spark usa lazy evaluation, no corren jobs de Spark
## hasta el momento
integersListRDD = sc.parallelize(integersList, 8)
type(integersListRDD)

### Asi se crean rdd con datos

datos = [
    (0, 1, 2, 2017, 20),
    (2, 2, 3, 2017, 30),
    (1, 1, 7, 2017, 15),
    (0, 1, 3, 2015, 23),
    (3, 3, 1, 2017, 10),
    (3, 3, 8, 2016, 17),
    (3, 3, 6, 2017, 5),
    (0, 4, 10, 2017, 12),
    (15, 2, 5, 2017, 15),
    (0, 5, 1, 2017, 70)
]

rdd_ejemplo = sc.parallelize(datos)
rdd_ejemplo.take(3)

[(0, 1, 2, 2017, 20), (2, 2, 3, 2017, 30), (1, 1, 7, 2017, 15)]

In [147]:
## podemos ver tambien otra informacion interesante del RDD
## el numero de particiones
integersListRDD.getNumPartitions()

8

In [148]:
## el conjunto de transformaciones que se aplica
integersListRDD.toDebugString()

b'(8) PythonRDD[176] at RDD at PythonRDD.scala:53 []\n |  ParallelCollectionRDD[172] at parallelize at PythonRDD.scala:195 []'

In [149]:
## para ver mas metodos disponibles del RDD
#help(integersListRDD)

In [150]:
integersListRDD.take(5)

[1, 2, 3, 4, 5]

In [151]:
integersListRDD.count()

1000

## Leyendo archivo con textFile

In [152]:
#downloaded = drive.CreateFile({'id':"1ybtSQxrqVqbRrl_3FMzMYW03Flp4zM-j"})   # replace the id with id of file you want to access
#downloaded.GetContentFile('s.txt') 
rdd = sc.textFile('shakespeare.txt')

In [153]:
rdd

shakespeare.txt MapPartitionsRDD[180] at textFile at NativeMethodAccessorImpl.java:0

In [154]:
rdd.count()

124614

In [155]:
rdd.take(5)

['1609', '', 'THE SONNETS', '', 'by William Shakespeare']

## Leyendo datos con el sqlContext

In [156]:
sqlContext = SQLContext(sc)

In [157]:
#dataframe = sqlContext.read.text('s.txt')

In [158]:
#dataframe

In [159]:
#rddCsv = dataframe.rdd

In [160]:
#rddCsv

In [161]:
#rddCsv.take(5)

Tambien se pueden leer archivos csv, json, parquet, jdbc, etc

# Acciones

## Count

Obtiene la cantidad de registros del RDD

In [162]:
integersListRDD.count()

1000

### Take

Obtiene los primeros n registros del RDD

In [163]:
integersListRDD.take(5)

[1, 2, 3, 4, 5]

## Collect

Obtiene TODOS los registros del RDD. Esto es un potencial problema, ya que si los datos no son acotados va a sobrecargar el driver. Solo se debe ejecutar si de antemano conocemos que la cantidad de datos es acotada.

In [164]:
integersListRDD.collect() # danger

[1,
 2,
 3,
 4,
 5,
 6,
 7,
 8,
 9,
 10,
 11,
 12,
 13,
 14,
 15,
 16,
 17,
 18,
 19,
 20,
 21,
 22,
 23,
 24,
 25,
 26,
 27,
 28,
 29,
 30,
 31,
 32,
 33,
 34,
 35,
 36,
 37,
 38,
 39,
 40,
 41,
 42,
 43,
 44,
 45,
 46,
 47,
 48,
 49,
 50,
 51,
 52,
 53,
 54,
 55,
 56,
 57,
 58,
 59,
 60,
 61,
 62,
 63,
 64,
 65,
 66,
 67,
 68,
 69,
 70,
 71,
 72,
 73,
 74,
 75,
 76,
 77,
 78,
 79,
 80,
 81,
 82,
 83,
 84,
 85,
 86,
 87,
 88,
 89,
 90,
 91,
 92,
 93,
 94,
 95,
 96,
 97,
 98,
 99,
 100,
 101,
 102,
 103,
 104,
 105,
 106,
 107,
 108,
 109,
 110,
 111,
 112,
 113,
 114,
 115,
 116,
 117,
 118,
 119,
 120,
 121,
 122,
 123,
 124,
 125,
 126,
 127,
 128,
 129,
 130,
 131,
 132,
 133,
 134,
 135,
 136,
 137,
 138,
 139,
 140,
 141,
 142,
 143,
 144,
 145,
 146,
 147,
 148,
 149,
 150,
 151,
 152,
 153,
 154,
 155,
 156,
 157,
 158,
 159,
 160,
 161,
 162,
 163,
 164,
 165,
 166,
 167,
 168,
 169,
 170,
 171,
 172,
 173,
 174,
 175,
 176,
 177,
 178,
 179,
 180,
 181,
 182,
 183,
 184,
 185

## First

Obtiene el primer registro del RDD

In [165]:
integersListRDD.first()

1

## TakeOrdered

Obtiene los primeros n registros en base a un orden indicado.

In [166]:
integersListRDD.takeOrdered(5, key=lambda x: -x)

[1000, 999, 998, 997, 996]

## TakeSample

Obtiene una muestra de n registros con o sin reemplazo.

In [167]:
integersListRDD.takeSample(False, 5)

[655, 941, 296, 502, 599]

## Reduce

Obtiene un solo registro, combinando el resultado en base a una función dada.

Suma de todos los nros del RDD:

In [168]:
integersListRDD.reduce(lambda a,b: a+b)

500500

Número más grande del RDD:

In [169]:
integersListRDD.reduce(lambda a,b: a if a > b else b)

1000

## CountByKey

Cuenta ocurrencias de registros para cada clave.

En Spark para que un registro sea considerado con clave debe se una tupla de unicamente dos elementos. El primer elemento es la key y el segundo el valor. A su vez, la key y el valor pueden estar compuestos por tuplas.

Cuento cuántos nros múltiplo de 2 hay y cuántos no:

In [170]:
integersListRDD.map(lambda x: (x % 2, 1)).countByKey()

defaultdict(int, {1: 500, 0: 500})

# Transformaciones

### Map

Transforma cada registro en base a la función dada.

In [171]:
integersListRDD.map(lambda x: x*2).take(5)

[2, 4, 6, 8, 10]

In [172]:
integersListRDD.map(lambda x: (x % 2, x)).take(5)

[(1, 1), (0, 2), (1, 3), (0, 4), (1, 5)]

## Filter

Filtra registros en base a la función dada.

In [173]:
integersListRDD.filter(lambda x: x % 2 == 0).take(5)

[2, 4, 6, 8, 10]

In [174]:
integersListRDD.filter(lambda x: x % 2 == 0).count()

500

## FlatMap

Similar a Map, pero cada registro puede generar 0, 1 o más registros.

Para cada registro original genero un nuevo registro con el nro, otro con el nro menos uno y otro con el registro más uno:

In [175]:
integersFlat = integersListRDD.flatMap(lambda x: [(x), (x-1), (x+1)])
integersFlat.count()

3000

## ReduceByKey

Combina los registros para una misma clave en base a una función de reduce.

La función de reduce debe ser **conmutativa** y **asociativa**.

Del RDD salida del flatMap cuento cuantos registros hay para cada nro:

In [176]:
integersFlat.map(lambda x: (x, 1)).reduceByKey(lambda a,b: a+b).count()

1002

In [177]:
integersFlat.map(lambda x: (x, 1)).reduceByKey(lambda a,b: a+b).take(10)

[(0, 1),
 (8, 3),
 (16, 3),
 (24, 3),
 (32, 3),
 (40, 3),
 (48, 3),
 (56, 3),
 (64, 3),
 (72, 3)]

In [178]:
integersFlat.map(lambda x: (x, 1)).reduceByKey(lambda a,b: a+b).reduce(lambda a,b: a if a > b else b)

(1001, 1)

## GroupByKey

Agrupa los registros para cada clave. Es similar a reduceByKey pero con 
groupByKey se obtiene todos los registros para cada clave.

Solo se debe utilizar si es necesario la información de cada registro y la cantidad de registros por clave no es demasiado grande.

GroupByKey es una transformación costosa.

Si se desea realizar una agregación, usar reduceByKey. Usar groupByKey para hacer una agregación esta MAL.

Necesito saber cuales son los nros múltiplos de 2 y cuales no:

In [179]:
integersListRDD.map(lambda x: (x % 2, x)).groupByKey().map(lambda x: (x[0], list(x[1]))).collect()

[(0,
  [2,
   4,
   6,
   8,
   10,
   12,
   14,
   16,
   18,
   20,
   22,
   24,
   26,
   28,
   30,
   32,
   34,
   36,
   38,
   40,
   42,
   44,
   46,
   48,
   50,
   52,
   54,
   56,
   58,
   60,
   62,
   64,
   66,
   68,
   70,
   72,
   74,
   76,
   78,
   80,
   82,
   84,
   86,
   88,
   90,
   92,
   94,
   96,
   98,
   100,
   102,
   104,
   106,
   108,
   110,
   112,
   114,
   116,
   118,
   120,
   122,
   124,
   126,
   128,
   130,
   132,
   134,
   136,
   138,
   140,
   142,
   144,
   146,
   148,
   150,
   152,
   154,
   156,
   158,
   160,
   162,
   164,
   166,
   168,
   170,
   172,
   174,
   176,
   178,
   180,
   182,
   184,
   186,
   188,
   190,
   192,
   194,
   196,
   198,
   200,
   202,
   204,
   206,
   208,
   210,
   212,
   214,
   216,
   218,
   220,
   222,
   224,
   226,
   228,
   230,
   232,
   234,
   236,
   238,
   240,
   242,
   244,
   246,
   248,
   250,
   252,
   254,
   256,
   258,
   260,
   262,


## Distinct

Elimina registros duplicados (todo el registro debe coincidir)

Del RDD trás aplicar flatMap, obtengo los registros únicos:

In [180]:
integersFlat.distinct().count()


1002

# Ejemplos de transformaciones y acciones con los textos de Shakespeare

## Leo de a líneas

In [181]:
lines = sc.textFile('shakespeare.txt')

## Cantidad de líneas totales

In [182]:
lines.count()

124614

## Primeras 10 líneas

In [183]:
lines.take(10)

['1609',
 '',
 'THE SONNETS',
 '',
 'by William Shakespeare',
 '',
 '',
 '',
 '                     1',
 '  From fairest creatures we desire increase,']

## Obtengo las palabras de todas las líneas (flatMap)

In [184]:
words = lines.flatMap(lambda x: x.split())

In [185]:
words.take(10)

['1609',
 'THE',
 'SONNETS',
 'by',
 'William',
 'Shakespeare',
 '1',
 'From',
 'fairest',
 'creatures']

In [186]:
words.count()

902892

## Contando palabras (reduceByKey)

In [187]:
wordsCount = words.map(lambda x: (x.lower(),1))

In [188]:
wordsCount.take(10)

[('1609', 1),
 ('the', 1),
 ('sonnets', 1),
 ('by', 1),
 ('william', 1),
 ('shakespeare', 1),
 ('1', 1),
 ('from', 1),
 ('fairest', 1),
 ('creatures', 1)]

In [189]:
wordsCounted = wordsCount.reduceByKey(lambda x,y: x+y)

In [190]:
wordsCounted.take(10)

[('shakespeare', 258),
 ('1', 13),
 ('fairest', 39),
 ('creatures', 27),
 ('we', 3210),
 ('increase,', 9),
 ('thereby', 21),
 ("beauty's", 30),
 ('rose', 44),
 ('never', 959)]

In [191]:
wordsCounted.takeOrdered(10, lambda x: -x[1])

[('the', 27681),
 ('and', 26066),
 ('i', 19540),
 ('to', 18737),
 ('of', 18084),
 ('a', 14424),
 ('my', 12456),
 ('in', 10721),
 ('you', 10666),
 ('that', 10489)]

### Mal uso de groupByKey

In [192]:
wordsCount.groupByKey().takeOrdered(10, lambda x: -1 * len(x[1]))

[('the', <pyspark.resultiterable.ResultIterable at 0x7f88d93fb850>),
 ('and', <pyspark.resultiterable.ResultIterable at 0x7f88d9420a10>),
 ('i', <pyspark.resultiterable.ResultIterable at 0x7f88d93fb750>),
 ('to', <pyspark.resultiterable.ResultIterable at 0x7f88d94204d0>),
 ('of', <pyspark.resultiterable.ResultIterable at 0x7f88d93fbb90>),
 ('a', <pyspark.resultiterable.ResultIterable at 0x7f88d94209d0>),
 ('my', <pyspark.resultiterable.ResultIterable at 0x7f88d9420c50>),
 ('in', <pyspark.resultiterable.ResultIterable at 0x7f88d93fbb10>),
 ('you', <pyspark.resultiterable.ResultIterable at 0x7f88d94206d0>),
 ('that', <pyspark.resultiterable.ResultIterable at 0x7f88d9420710>)]

In [193]:
wordsCount.groupByKey().map(lambda a: (a[0], list(a[1]))).take(5)

[('shakespeare',
  [1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,
   1,

In [194]:
wordsCount.groupByKey().filter(lambda x: len(x[1]) < 5).map(lambda a: (a[0], list(a[1]))).take(5)

[('riper', [1, 1, 1]),
 ('memory:', [1]),
 ("feed'st", [1, 1, 1]),
 ("light's", [1]),
 ('fuel,', [1])]

## Palabra más larga (reduce)

In [195]:
words.reduce(lambda a, b: a if (len(a) > len(b)) else b)

'http://www.ibiblio.org/gutenberg/etext06'

## Palabras que empiezan con a (filter)

In [196]:
wordsA = words.filter(lambda word: word.startswith('a'))

In [197]:
wordsA.count()

63676

In [198]:
wordsA.take(10)

['as', 'a', 'abundance', 'art', 'and', 'a', 'asked,', 'all', 'all', 'an']

## Palabras únicas que empiezan con a (distinct)

In [199]:
wordsA.distinct().count()

2688

## Cantidad de palabras por frecuencia de repetición ordenados (sortByKey)

In [200]:
wordsCounted.take(5)

[('shakespeare', 258),
 ('1', 13),
 ('fairest', 39),
 ('creatures', 27),
 ('we', 3210)]

In [201]:
wordsFreq = wordsCounted.map(lambda x: (x[1],1))

In [202]:
wordsFreq.take(10)

[(258, 1),
 (13, 1),
 (39, 1),
 (27, 1),
 (3210, 1),
 (9, 1),
 (21, 1),
 (30, 1),
 (44, 1),
 (959, 1)]

In [203]:
wordsFreq.reduceByKey(lambda a,b: a+b).sortByKey().take(10)

[(1, 31072),
 (2, 8493),
 (3, 4342),
 (4, 2659),
 (5, 1822),
 (6, 1338),
 (7, 1053),
 (8, 779),
 (9, 700),
 (10, 549)]

In [204]:
wordsFreq.reduceByKey(lambda a,b: a+b).takeOrdered(10, lambda x: -x[1])

[(1, 31072),
 (2, 8493),
 (3, 4342),
 (4, 2659),
 (5, 1822),
 (6, 1338),
 (7, 1053),
 (8, 779),
 (9, 700),
 (10, 549)]

# Transformaciones entre dos RDD

## Union

Obtiene la unión entre dos RDD.

In [205]:
integersList2 = range(501,1501)
len(integersList2)

1000

In [206]:
integersList2RDD = sc.parallelize(integersList2)

In [207]:
integersList2RDD.count()

1000

In [208]:
integersListRDD.count()

1000

In [209]:
union = integersListRDD.union(integersList2RDD)

In [210]:
union.take(5)

[1, 2, 3, 4, 5]

In [211]:
union.count()

2000

## Intersection

Intersección entre dos RDD.

In [212]:
intersection = integersListRDD.intersection(integersList2RDD)

In [213]:
intersection.count()

500

In [214]:
intersection.take(10)

[512, 528, 544, 560, 576, 592, 608, 624, 640, 656]

In [215]:
intersection.collect()

[512,
 528,
 544,
 560,
 576,
 592,
 608,
 624,
 640,
 656,
 672,
 688,
 704,
 720,
 736,
 752,
 768,
 784,
 800,
 816,
 832,
 848,
 864,
 880,
 896,
 912,
 928,
 944,
 960,
 976,
 992,
 513,
 529,
 545,
 561,
 577,
 593,
 609,
 625,
 641,
 657,
 673,
 689,
 705,
 721,
 737,
 753,
 769,
 785,
 801,
 817,
 833,
 849,
 865,
 881,
 897,
 913,
 929,
 945,
 961,
 977,
 993,
 514,
 530,
 546,
 562,
 578,
 594,
 610,
 626,
 642,
 658,
 674,
 690,
 706,
 722,
 738,
 754,
 770,
 786,
 802,
 818,
 834,
 850,
 866,
 882,
 898,
 914,
 930,
 946,
 962,
 978,
 994,
 515,
 531,
 547,
 563,
 579,
 595,
 611,
 627,
 643,
 659,
 675,
 691,
 707,
 723,
 739,
 755,
 771,
 787,
 803,
 819,
 835,
 851,
 867,
 883,
 899,
 915,
 931,
 947,
 963,
 979,
 995,
 516,
 532,
 548,
 564,
 580,
 596,
 612,
 628,
 644,
 660,
 676,
 692,
 708,
 724,
 740,
 756,
 772,
 788,
 804,
 820,
 836,
 852,
 868,
 884,
 900,
 916,
 932,
 948,
 964,
 980,
 996,
 501,
 517,
 533,
 549,
 565,
 581,
 597,
 613,
 629,
 645,
 661,
 677

## Subtract

Elimina del primer RDD los registros que aparezcan en el segundo.

In [216]:
subtract = integersListRDD.subtract(integersList2RDD)

In [217]:
subtract.count()

500

In [218]:
subtract.collect()

[16,
 32,
 48,
 64,
 80,
 96,
 112,
 128,
 144,
 160,
 176,
 192,
 208,
 224,
 240,
 256,
 272,
 288,
 304,
 320,
 336,
 352,
 368,
 384,
 400,
 416,
 432,
 448,
 464,
 480,
 496,
 1,
 17,
 33,
 49,
 65,
 81,
 97,
 113,
 129,
 145,
 161,
 177,
 193,
 209,
 225,
 241,
 257,
 273,
 289,
 305,
 321,
 337,
 353,
 369,
 385,
 401,
 417,
 433,
 449,
 465,
 481,
 497,
 2,
 18,
 34,
 50,
 66,
 82,
 98,
 114,
 130,
 146,
 162,
 178,
 194,
 210,
 226,
 242,
 258,
 274,
 290,
 306,
 322,
 338,
 354,
 370,
 386,
 402,
 418,
 434,
 450,
 466,
 482,
 498,
 3,
 19,
 35,
 51,
 67,
 83,
 99,
 115,
 131,
 147,
 163,
 179,
 195,
 211,
 227,
 243,
 259,
 275,
 291,
 307,
 323,
 339,
 355,
 371,
 387,
 403,
 419,
 435,
 451,
 467,
 483,
 499,
 4,
 20,
 36,
 52,
 68,
 84,
 100,
 116,
 132,
 148,
 164,
 180,
 196,
 212,
 228,
 244,
 260,
 276,
 292,
 308,
 324,
 340,
 356,
 372,
 388,
 404,
 420,
 436,
 452,
 468,
 484,
 500,
 5,
 21,
 37,
 53,
 69,
 85,
 101,
 117,
 133,
 149,
 165,
 181,
 197,
 213,
 229,


## Joins

Con los joins se combinan dos RDD en base a las claves de los registros. Junta cada registro del primer RDD con cada registro del segundo RDD que tengan la misma clave. No agrupa, sino que es de a pares de registro.

In [219]:
data_alumnos = [
  (1,'Damian'),
  (2,'Luis'),
  (3,'Martin'),
  (4,'Natalia'),
  (5,'Joaquin')
]

alumnos = sc.parallelize(data_alumnos)

In [220]:
alumnos.collect()

[(1, 'Damian'), (2, 'Luis'), (3, 'Martin'), (4, 'Natalia'), (5, 'Joaquin')]

In [221]:
data_materias_aprobadas = [
  (1, 'Algebra'),
  (2, 'Análisis Matemático'),
  (200, 'Algebra'),
  (2, 'Física')
]

materias_aprobadas = sc.parallelize(data_materias_aprobadas)

In [222]:
materias_aprobadas.collect()

[(1, 'Algebra'), (2, 'Análisis Matemático'), (200, 'Algebra'), (2, 'Física')]

### Inner Join (Join)

Cuando se llama para sets de datos del tipo (K,V) y (K,W) devuelve un set de datos del tipo (K, (V,W)) con todos los pares de elementos para cada key. (especificamente los que hay en comun por esa clave en ambos sets de datos)

In [223]:
alumnos.join(materias_aprobadas).collect()

[(1, ('Damian', 'Algebra')),
 (2, ('Luis', 'Análisis Matemático')),
 (2, ('Luis', 'Física'))]

### Left Outer Join

Cuando se llama para sets de datos del tipo (K,V) y (K,W) devuelve un set de datos del tipo (K, (V,W)) asegurandonos que todos los datos del set de datos izquierdo estaran en el resultado del join.

In [224]:
alumnos.leftOuterJoin(materias_aprobadas).collect()

[(1, ('Damian', 'Algebra')),
 (2, ('Luis', 'Análisis Matemático')),
 (2, ('Luis', 'Física')),
 (3, ('Martin', None)),
 (4, ('Natalia', None)),
 (5, ('Joaquin', None))]

### Right Outer Join

Cuando se llama para sets de datos del tipo (K,V) y (K,W) devuelve un set de datos del tipo (K, (V,W)) asegurandonos que todos los datos del set de datos derecho estaran en el resultado del join.

In [225]:
alumnos.rightOuterJoin(materias_aprobadas).collect()

[(1, ('Damian', 'Algebra')),
 (2, ('Luis', 'Análisis Matemático')),
 (2, ('Luis', 'Física')),
 (200, (None, 'Algebra'))]

### Outer/Full Join


Cuando se llama para sets de datos del tipo (K,V) y (K,W) devuelve un set de datos del tipo (K, (V,W)) asegurandonos que todos los datos de ambos set de datos estaran aunque no haya match de keys.

In [226]:
alumnos.fullOuterJoin(materias_aprobadas).collect()

[(1, ('Damian', 'Algebra')),
 (2, ('Luis', 'Análisis Matemático')),
 (2, ('Luis', 'Física')),
 (3, ('Martin', None)),
 (4, ('Natalia', None)),
 (5, ('Joaquin', None)),
 (200, (None, 'Algebra'))]

### Broadcast Join (map-side join)

#### Variable Broadcast

Una variable Broadcast nos permite mantener una variable solo lectura cacheada en cada una de las maquinas del cluster en vez de enviar esa informacion con cada una de las tareas que se envian al cluster.

Esto es particularmente util cuando cuando tareas a partir de multiples etapas (stages) necesitan la misma información o cuando cachear información de forma deserializada es importante.

Tener en cuenta que esto **es posible** cuando uno de los data sets o conjunto de datos **es lo suficientemente pequeño para ser broadcasteado a todos los nodos/workers del cluster**.

In [227]:
# Vamos a suponer que tenemos un RDD de productos por sus IDs identificando ventas de los mismos
prodsList = [1,11,1,4,5,11,2,3,4,5,6,4,5,4,3,2,1,11,2,3,4,5,6,4,3,2,1,1]
prods = sc.parallelize(prodsList,3)

In [228]:
# Un hash con los productos y sus nombres
productNames = {1:'papas',
                2:'cebollas',
                3:'tomates',
                4:'zanahorias',
                5:'batatas',
                6:'peras',
                7:'cilantro',
                8:'apio',
                9:'morrones',
                10:'manzanas',
                11:'naranjas'}

# Hacemos un broadcast de la variable
bproductNames = sc.broadcast(productNames)

In [229]:
# Buscamos los productos que se vendieron más de 4 veces
popularProds = prods.map(lambda x:(x,1))\
    .reduceByKey(lambda x,y:x+y)\
    .filter(lambda x:x[1]>=4)
popularProds.collect()

[(3, 4), (1, 5), (4, 6), (5, 4), (2, 4)]

El join se realiza de forma implicita usando un map y dentro del mismo accediendo a la informacion de la variable a la que se realizo el broadcast via .value

In [230]:
popularProds = popularProds.map(
    lambda x:(bproductNames.value[x[0]],x[0],x[1]))
popularProds.collect()

[('tomates', 3, 4),
 ('papas', 1, 5),
 ('zanahorias', 4, 6),
 ('batatas', 5, 4),
 ('cebollas', 2, 4)]

#### Ventajas

Cuando un valor es "broadcasteado" al cluster, este es copiado a los nodos/workers **sólo una vez** (en vez de múltiples veces si la información fuera a enviarse en cada task). De esta forma se resuelve la consulta más rapidamente.

# Transformaciones sobre las particiones

In [231]:
rdd = sc.parallelize(range(1,11))
rdd.getNumPartitions()

8

In [232]:
sc.defaultParallelism

8

In [233]:
rdd.collect()

[1, 2, 3, 4, 5, 6, 7, 8, 9, 10]

## Glom

Junta los registros de cada partición en una lista.

In [234]:
rdd.glom().collect()

[[1], [2], [3], [4, 5], [6], [7], [8], [9, 10]]

## MapPartitions

Devuelve un nuevo RDD aplicando una función a cada partición del RDD.

In [235]:
def f(iterator): yield __builtin__.sum(iterator)
rdd.mapPartitions(f).collect()

[1, 2, 3, 9, 6, 7, 8, 19]

## Repartition

Reshuffle los datos en el RDD de forma aleatoria para crear más o menos particiones y balancearlas. 

Hace un shuffle de todo los datos por la red.

In [236]:
rdd = sc.parallelize(range(1,11), 4)
rdd.getNumPartitions()

4

In [237]:
rdd.glom().collect()

[[1, 2], [3, 4, 5], [6, 7], [8, 9, 10]]

In [238]:
rdd2 = rdd.repartition(2)

In [239]:
rdd2.getNumPartitions()

2

In [240]:
rdd2.glom().collect()

[[1, 2, 6, 7, 8, 9, 10], [3, 4, 5]]

Spark no hace shuffle de registros individuales sino de a bloques con un mínimo (no es un problema cuando se manejan grandes cantidades de datos)

## Coalesce

Decrementa la cantidad de particiones del RDD.

No hace shuffle por defecto, solo pasa datos de una partición a otra.

No quedan balanceadas.

In [241]:
rddCoalesce = rdd.coalesce(2)

In [242]:
rddCoalesce.glom().collect()

[[1, 2, 3, 4, 5], [6, 7, 8, 9, 10]]

## RepartitionAndSortWithinPartitions

Reparticiona un RDD de acuerdo a un particionador y ordena los registros en base a su clave.

Los registros deben tener clave.

Es más eficiente que hacer un repartition y luego un sort dentro de cada partición ya que realiza el sort en el mismo paso de shuffle.

In [243]:
rdd.map(lambda x: (x, x)).collect()

[(1, 1),
 (2, 2),
 (3, 3),
 (4, 4),
 (5, 5),
 (6, 6),
 (7, 7),
 (8, 8),
 (9, 9),
 (10, 10)]

In [244]:
rdd.map(lambda x: (x, x)).glom().collect()

[[(1, 1), (2, 2)],
 [(3, 3), (4, 4), (5, 5)],
 [(6, 6), (7, 7)],
 [(8, 8), (9, 9), (10, 10)]]

### Ascending

In [245]:
rdd.map(lambda x: (x, x)).repartitionAndSortWithinPartitions(2).glom().collect()

[[(2, 2), (4, 4), (6, 6), (8, 8), (10, 10)],
 [(1, 1), (3, 3), (5, 5), (7, 7), (9, 9)]]

In [246]:
rdd.map(lambda x: (x % 3, x)).repartitionAndSortWithinPartitions(2).glom().collect()

[[(0, 3), (0, 6), (0, 9), (2, 2), (2, 5), (2, 8)],
 [(1, 1), (1, 4), (1, 7), (1, 10)]]

In [247]:
rdd.map(lambda x: (x % 3, x)).repartitionAndSortWithinPartitions(2, ascending=False).glom().collect()

[[(2, 2), (2, 5), (2, 8), (0, 3), (0, 6), (0, 9)],
 [(1, 1), (1, 4), (1, 7), (1, 10)]]

### PartitionFunc

In [248]:
rdd.map(lambda x: (x * 2, x)).repartitionAndSortWithinPartitions(2).glom().collect()

[[(2, 1),
  (4, 2),
  (6, 3),
  (8, 4),
  (10, 5),
  (12, 6),
  (14, 7),
  (16, 8),
  (18, 9),
  (20, 10)],
 []]

In [249]:
rdd.map(lambda x: (x * 2, x)).repartitionAndSortWithinPartitions(2, partitionFunc=lambda x: (x % 3)).glom().collect()

[[(2, 1), (6, 3), (8, 4), (12, 6), (14, 7), (18, 9), (20, 10)],
 [(4, 2), (10, 5), (16, 8)]]

# Persistiendo RDD

## Cache

Cachea un RDD intermedio que va a ser utilizado varias veces de modo de evitar tener que ejecutar todas las transformaciones cada vez.

In [250]:
rdd = sc.parallelize(range(1,100000))

In [251]:
rddCached = rdd.map(lambda x: x*10).cache()

In [252]:
rddCached.count()

99999

In [253]:
rddCached.take(10)

[10, 20, 30, 40, 50, 60, 70, 80, 90, 100]

## SaveAsTextFile

Guarda un RDD a disco en un archivo de texto.

In [254]:
rdd.saveAsTextFile('numbers.txt')

In [255]:
rddN = sc.textFile('numbers.txt')

In [256]:
rddN.collect()

['62500',
 '62501',
 '62502',
 '62503',
 '62504',
 '62505',
 '62506',
 '62507',
 '62508',
 '62509',
 '62510',
 '62511',
 '62512',
 '62513',
 '62514',
 '62515',
 '62516',
 '62517',
 '62518',
 '62519',
 '62520',
 '62521',
 '62522',
 '62523',
 '62524',
 '62525',
 '62526',
 '62527',
 '62528',
 '62529',
 '62530',
 '62531',
 '62532',
 '62533',
 '62534',
 '62535',
 '62536',
 '62537',
 '62538',
 '62539',
 '62540',
 '62541',
 '62542',
 '62543',
 '62544',
 '62545',
 '62546',
 '62547',
 '62548',
 '62549',
 '62550',
 '62551',
 '62552',
 '62553',
 '62554',
 '62555',
 '62556',
 '62557',
 '62558',
 '62559',
 '62560',
 '62561',
 '62562',
 '62563',
 '62564',
 '62565',
 '62566',
 '62567',
 '62568',
 '62569',
 '62570',
 '62571',
 '62572',
 '62573',
 '62574',
 '62575',
 '62576',
 '62577',
 '62578',
 '62579',
 '62580',
 '62581',
 '62582',
 '62583',
 '62584',
 '62585',
 '62586',
 '62587',
 '62588',
 '62589',
 '62590',
 '62591',
 '62592',
 '62593',
 '62594',
 '62595',
 '62596',
 '62597',
 '62598',
 '62599',


## SaveAsPickleFile

Guarda un RDD a disco en un archivo con los datos serializados.

In [257]:
rdd.saveAsPickleFile('numbers2.file')

In [258]:
rddN2 = sc.pickleFile('numbers2.file')

In [259]:
rddN2.take(10)

[62500, 62501, 62502, 62503, 62504, 62505, 62506, 62507, 62508, 62509]