martes, mayo 20, 2014

Usando Amazon EMR para ejecutar un Hadoop MapReduce

AWS Elastic MapReduce (EMR)

Lo que vamos a hacer es usar las capacidad de Amazon EMR (Elastic Map Reduce) para ejecutar nuestro contador de palabras sin tener que preocuparnos de instalaciones de clusters.

NOTA IMPORTANTE:

Amazon EMR NO es gratuito y no entra dentro de las posibilidades de la AWS Free Tier por lo que realizar este tutorial, aunque cuesta menos de 1€, no es gratuito. Al final del todo pondré una captura de pantalla del coste total que fueron 0.21€

Creando un bucket S3 y subiendo los archivos

Lo primero que tendremos que hacer es crear un bucket S3 para subir archivos dentro. S3 es un sistema de almacenamiento en la nube que, con la AWS Free tier nos otorga 5gb de almacenamiento gratuito. Habrá que meterse en la cuenta de AWS y seleccionar, dentro de los servicios, el S3:


Una vez en la ventana de S3, hay que hacer click en "Create bucket":

Y darle un nombre al bucket:

Subir los archivos al bucket

Ahora que ya tenemos un bucket de almacenamiento, hay que seleccionar "Actions-Upload" para seleccionar los dos archivos que queremos subir, uno de ellos el fichero del que queremos contar las palabras y otro el JAR que contiene el Wordcount que escribimos en el tutorial de 



Para este caso, he subido una versión del Quijote en texto plano que encontré por la red y el JAR con el Wordcount . El bucket debe presentar este aspecto:

Creando el cluster

Paso 1: Cluster configuration

Ahora, seleccionamos el servicio de Amazon EMR de la selección de servicios. En la ventana que nos aparece hay que seleccionar "Create cluster":

Primer paso para la creación de un cluster



Atención al botón de "Configure sample application". Como se puede intuir, este botón nos abre otra ventana para poder ejecutar el cluster con algunos trabajos de ejemplo (entre ellos el propio contador de palabras). Como lo que queremos es aprender como se hace, vamos a evitar usar el botón.

Paso 2: Software configuration

Una vez le hemos dado un nombre al cluster y hemos seleccionado una ubicación en nuestro S3 para el almacenamiento de los logs, pasamos a seleccionar la AMI y versión de Hadoop a usar:
Vamos a usar la version 2.4.2 de la AMI con hadoop 1.0.3. Como podeis ver, también nos da la opción de instalar Hive o Pig. Para los propósitos de este tutorial no será necesario pero no pasa nada si se instalan.

Paso 3: Hardware configuration

Pasamos al siguiente paso que va a ser la configuración de acceso. Para los propósitos de este tutorial vamos a usar una instancia small para el contenedor del Namenode, el SecondaryNamenode y el JobTracker y 2 instancias small con un DataNode y un TaskTracker cada una. Por último, aunque no es estrictamente necesario, hay que seleccionar una clave PEM para el acceso a las máquinas. A nosotros no nos va a hacer falta ya que sólo vamos a arrancar las máquinas para hacer el conteo y se van a parar al finalizar.

Paso 4: Steps

El último de los pasos para la creación del cluster son los "Steps", que son los trabajos a realizar. Nuestro step va a ser ejecutar el Wordcount que vendrá configurado automáticamente (lo seleccionamos en un paso anterior cuando abrimos el cuadro de diálogo de "Configure Sample Application").


Debemos editar el paso para indicarle el fichero del que queremos contar el número de palabras. En la nueva ventana que aparece, sólo tendremos que abrir el cuadro de diálogo de "Input S3 Location" y seleccionar el fichero de "quijote.txt" como se puede apreciar a continuación:
Asegurarse de tener elegido "Terminate cluster" en "Action on Failure"

Navegar y seleccionar el archivo a contar palabras

Una vez tengamos todo, seleccionamos "Select" dentro de la ventana de "Select S3 Folder" y "Save" en la ventana de "Add Step". Ya está el trabajo y el cluster configurado para ejecutar.

Finalmente hacemos click en "Create cluster" y nuestro MapReduce ya estará corriendo.

Algunos aspectos básicos de AWS EMR: Si por alguna circunstancia hay un error de configuración dentro del cluster, no se cargará ningún gasto a vuestra cuenta. Sólo se carga cuando el MapReduce consigue iniciarse, en cuyo caso, si hay error; sí que habrá un pequeño cargo.

Resultado


Los resultados de la ejecución los podremos ver en nuestro S3, dentro de la carpeta que venía configurada  en el "Output S3 folder" dentro de la ventana de "Add Step". En la siguiente imagen podéis ver en el resultado del Wordcount dentro del bucket S3:


Notas finales

Como se puede apreciar, la ejecución de una tarea EMR en AWS es realmente sencilla, mucho mas que crear nuestro propio cluster, nuestro propio MapReduce, etc. Esto demuestra las tremendas posibilidades de usar AWS EMR para la ejecución de docenas o cientos de "Workers" a sólo unos clicks de ratón. A mi, personalmente, me dejó impresionado, sobre todo teniendo en cuenta que también se pueden ejecutar scripts Pig y Hive con la misma facilidad.

lunes, mayo 19, 2014

Instalando Hadoop en modo distribuido

Hadoop "Fully Distributed"

Aunque pueda asustar, la instalación de un sistema distribuido de Hadoop no es tan complicado como pueda parecer. En realidad, siguiendo el tutorial sobre como instalarlo en modo pseudo-distribuido ya tenemos la mayor parte del trabajo hecho: Instalando Hadoop en modo pseudo-distribuido

Las diferencias  mas simples entre el modo distribuido y el modo pseudo distribuido son las siguientes:
  1. El nodo NameNode debe poder acceder por SSH sin contraseña a todos los nodos esclavos y a si mismo. (no es necesario que los nodos esclavos puedan acceder al nodo maestro).
  2. Todos los nodos deben compartir la misma configuración de Hadoop (el contenido de "conf", "etc/conf", o "/etc/hadoop" dependiendo de la versión).
  3. La reglas de Firewall deben permitir la comunicación via TCP entre todos los tipos de nodos en modo bidireccional

Comenzando

Lo primero, vamos a usar dos portátiles que vamos a llamar E1 y E2. Tenemos que instalar Hadoop en los dos equipos como si fuera modo pseudo-distribuido usando el tutorial del link anterior. Ahora lo importante es cambiar las direcciones IP's por direcciones reales dentro de nuestra LAN por lo que, si nuestro localhost es el 192.168.1.10 dentro de nuestra LAN, tendremos que usar esta IP.

Una vez tengamos los dos equipos con Hadoop en modo pseudo-distribuido vamos a ver como dividimos las tareas entre los dos equipos:
  1. E1: El equipo 1 va a ser el nodo "Maestro", por tanto vamos a tener en ejecución los siguientes hilos:
    1. Namenode
    2. SecondaryNamenode
    3. JobTracker
  2. E3: El Equipo 2 va a ser el nodo "Worker" por lo que las tareas en ejecución serán las siguientes:
    1. TaskTracker
    2. DataNode
Entraremos a la consola del E1, nuestra IP es 192.168.1.10. Vamos a editar, dentro de la carpeta de Hadoop, el fichero de configuración de esclavos donde vamos a indicar la IP del E2 (192.168.1.20). Dicho fichero lo encontraremos en conf/slaves:

192.168.1.20
Además, también hay que indicar la máquina que va a tener el SecondaryNamenode, para ello tenermos que editar el fichero conf/masters para agregrar la nuestra IP (la IP del E1):

192.168.1.10

NOTA: Aunque el fichero se llame "masters", no indica las máquinas donde se van a alojar el NameNode ni el Jobtracker, quedando estas configuraciones a disposición del fichero mapred-site.xml y del nodo que arranque el NameNode (en nuestro caso será el E1)

Configurar el acceso SSH

El E1 debe tener acceso por SSH sin contraseña al E2 y a sí mismo. Si aún no sabes como conseguir esto, está todo explicado en el tutorial para instalar Hadoop en modo pseudo-distribuido.

Configurar el mapred-site y el hdfs-site

Ahora, tendremos que tener comunes los ficheros hdfs-site y mapred-site a lo largo de nuestro cluster. En mi experiencia he encontrado bastante sencilla la sincronización a través de un repositorio git como puede ser Github pero cualquier solución es igualmente válida:







	
		dfs.replication
		2
	







    mapred.job.tracker
    192.168.1.10:54311




En el primero, el hdfs, estamos diciéndole que cada bloque de información lo duplique a lo largo del cluster (dfs-replication=2). En el segundo le estamos pasando la IP del E1 que es el que va a ejecutar el JobTracker.

Arrancando el cluster

Por último, hay que arrancar el cluster, desde el E1 y el directorio de instalación de Hadoop ejecutaremos:

NameNode
JobTracker
SecondaryNameNode

Y, si nos vamos al E2, podremos ver los siguientes:
DataNode
TaskTracker

Si resulta que te encuentras el Datanode y el Tasktracker dentro de la E1 es probablemente porque todavía conservas la IP de1 dentro del fichero conf/slaves.

Siguientes pasos, ¡instalar Hadoop en todas las máquinas que tenga por casa!

viernes, mayo 09, 2014

Escribiendo un Hadoop MapReduce en Java: Un Wordcount mejorado

En este tutorial, vamos a crear nuestro propio MapReduce contador de palabras y luego lo vamos a usar en nuestro cluster pseudo-distribuido de Hadoop.

Un trabajo Map Reduce se divide en 2 clases (Mapper y Reducer) mas el método de entrada "main". En este caso, el mapper va a crear un listado del tipo y el reducer va a "reducir" el listado de valores a un único valor que, en nuestro caso, será el número de ocurrencias de la clave.

/**
 *  Mapper.
 *  Recibirá, línea tras línea (el objeto Text), los contenidos del fichero que
 *  le pasemos por los argumentos.
 *
 *  Vamos a mejorar el Mapper inicial quitando todos los caracteres no alfabéticos
 *  de la cadena de entrada para evitar que, por ejemplo, "casa" y "casa." se
 *  contabilicen como palabras distintas
 */
public static class WordCountMapper extends Mapper {
  private Text word = new Text();

  @Override
  public void map(LongWritable key, Text value, Context context)
      throws IOException, InterruptedException 
  {
    //Escribimos una expresión regular para eliminar todo caracter no alfabético
    String cleanString = value.toString().replaceAll("[^A-Za-z\\s]", "");

    //Ahora cogemos el texto limpio y lo separamos por espacios
    StringTokenizer token = new StringTokenizer(cleanString);
    
    //Recorremos el tokenizer para recoger todas las palabras hasta que no haya mas
    while(token.hasMoreTokens())
    {
      //Palabra actual
      String tok = token.nextToken();

      //Asignar la palabra al objeto que se va a pasar al reducer
      word.set(tok);

      /*
       * Guardar la palabra con un valor número de 1 de tal manera que, si tenemos
       * la palabra "casa" estamos guardando  Si la palabra casa volviera
       * a aparecer, como ya la hemos establecido una vez, el resultado sería 
       *  y así por cada ocurrencia
       */
      context.write(word, new IntWritable(1));
    }
  }
}
Vamos ahora con el Reducer, ojo al atributo estático de la clase
/*
 * El reducer va a contar el número de ítems en la lista de palabras 
 * que pasamos desde el Mapper
 */
public static class WordCountReducer extends Reducer {

  @Override
  public void reduce(Text word, Iterable list, Context context)
      throws IOException, InterruptedException
  {
    //Ponemos el contador de palabras a 0
    int total = 0;
    
    //Recorremos el objeto Iterable list y sumamos uno al contador
    for(IntWritable count : list)
    {
      total++;
    }
    
    /* Escribimos el resultado del reducer, para el ejemplo de casa, escribiría
     * 
     */
    context.write(word, new IntWritable(total));
  }
  
}
Ahora solo falta la clase main de entrada a la aplicación donde se configura todas las opciones del MapReduce
// Entrada de la aplicación
public static void main(String[] args) throws IOException, ClassNotFoundException, InterruptedException 
{
  //Creamos un fichero de configuracion y le damos un nombre al Job
  Configuration conf = new Configuration();
  Job job = new Job(conf, "word count");

  //Le tenemos que indicar la clase que hay que usar para llamar a los mappers y reducers
  job.setJarByClass(Wordcount.class);

  //Le indicamos el nombre de la clase Mapper y de la clase Reducer
  job.setMapperClass(WordCountMapper.class);
  job.setReducerClass(WordCountReducer.class);

  /* 
   * Le tenemos que indicar el formato que va a tener el resultado, en nuestro caso vamos
   * a recuperar un resultado del tipo  por lo que, usando
   * los tipos primitivos de Hadoop, esto equivaldría a Text.class y IntWritable.class
   */
  job.setOutputKeyClass(Text.class);
  job.setOutputValueClass(IntWritable.class);

  /* Le indicamos el fichero de entrada y de salida que, por lo general, los vamos a recoger
   * de los parámetros que le pasemos a la clase JAR
   */
  FileInputFormat.addInputPath(job, new Path(args[0]));
  FileOutputFormat.setOutputPath(job, new Path(args[1]));

  //Esperamos a que el trabajo termine
  System.exit(job.waitForCompletion(true) ? 0 : 1);
}
Nos queda exportar el JAR y ejecutarlo con algún fichero de prueba para contar sus palabras. Para exportarlo vamos a usar Eclipse, por su amplia uso actual y su simplicidad. Tendremos que exportarlo como un JAR ejecutable asegurándonos de tener las siguientes opciones seleccionadas. Es importante seleccionar "Extract required libraries into generated JAR" en el cuadro de diálogo de exportación que nos aparece:

Ahora sólo falta probarlo. Para ello podemos usar cualquiera de los 3 modos de Hadoop (local, pseudo-distribuido o distribuido). Por simplicidad vamos a usar el modo Local:

cd $HADOOP_INSTALL
bin/hadoop jar mariocaster.blogspot.com-wordcount.jar example-files/input example-files/output
Recordar que la carpeta example-files/input debe contener los ficheros de texto a contar y la carpeta example-files/output NO debe existir cuando ejecutemos el script porque fallará (Hadoop se encargará de crear la carpeta). El resultado es similar al siguiente:

INFO util.NativeCodeLoader: Loaded the native-hadoop library
WARN mapred.JobClient: Use GenericOptionsParser for parsing the arguments. Applications should implement Tool for the same.
WARN mapred.JobClient: No job jar file set.  User classes may not be found. See JobConf(Class) or JobConf#setJar(String).
INFO input.FileInputFormat: Total input paths to process : 2
WARN snappy.LoadSnappy: Snappy native library not loaded
INFO mapred.JobClient: Running job: job_local1360868529_0001
INFO mapred.LocalJobRunner: Waiting for map tasks
INFO mapred.LocalJobRunner: Starting task: attempt_local1360868529_0001_m_000000_0
INFO util.ProcessTree: setsid exited with exit code 0
INFO mapred.Task:  Using ResourceCalculatorPlugin : org.apache.hadoop.util.LinuxResourceCalculatorPlugin@4f69385e
INFO mapred.MapTask: Processing split: file:/var/hadoop/examples-files/input/text2:0+118
INFO mapred.MapTask: io.sort.mb = 100
INFO mapred.MapTask: data buffer = 79691776/99614720
INFO mapred.MapTask: record buffer = 262144/327680
INFO mapred.MapTask: Starting flush of map output
INFO mapred.MapTask: Finished spill 0
INFO mapred.Task: Task:attempt_local1360868529_0001_m_000000_0 is done. And is in the process of commiting
INFO mapred.LocalJobRunner: 
INFO mapred.Task: Task 'attempt_local1360868529_0001_m_000000_0' done.
INFO mapred.LocalJobRunner: Finishing task: attempt_local1360868529_0001_m_000000_0
INFO mapred.LocalJobRunner: Starting task: attempt_local1360868529_0001_m_000001_0
INFO mapred.Task:  Using ResourceCalculatorPlugin : org.apache.hadoop.util.LinuxResourceCalculatorPlugin@4e3600d
INFO mapred.MapTask: Processing split: file:/var/hadoop/examples-files/input/texto:0+72
INFO mapred.MapTask: io.sort.mb = 100
INFO mapred.MapTask: data buffer = 79691776/99614720
INFO mapred.MapTask: record buffer = 262144/327680
INFO mapred.MapTask: Starting flush of map output
INFO mapred.MapTask: Finished spill 0
INFO mapred.Task: Task:attempt_local1360868529_0001_m_000001_0 is done. And is in the process of commiting
INFO mapred.LocalJobRunner: 
INFO mapred.Task: Task 'attempt_local1360868529_0001_m_000001_0' done.
INFO mapred.LocalJobRunner: Finishing task: attempt_local1360868529_0001_m_000001_0
INFO mapred.LocalJobRunner: Map task executor complete.
INFO mapred.Task:  Using ResourceCalculatorPlugin : org.apache.hadoop.util.LinuxResourceCalculatorPlugin@58017a75
INFO mapred.LocalJobRunner: 
INFO mapred.Merger: Merging 2 sorted segments
INFO mapred.Merger: Down to the last merge-pass, with 2 segments left of total size: 428 bytes
INFO mapred.LocalJobRunner: 
INFO mapred.Task: Task:attempt_local1360868529_0001_r_000000_0 is done. And is in the process of commiting
INFO mapred.LocalJobRunner: 
INFO mapred.Task: Task attempt_local1360868529_0001_r_000000_0 is allowed to commit now
INFO output.FileOutputCommitter: Saved output of task 'attempt_local1360868529_0001_r_000000_0' to examples-files/ouput
INFO mapred.LocalJobRunner: reduce > reduce
INFO mapred.Task: Task 'attempt_local1360868529_0001_r_000000_0' done.
INFO mapred.JobClient:  map 100% reduce 100%
INFO mapred.JobClient: Job complete: job_local1360868529_0001
INFO mapred.JobClient: Counters: 20
INFO mapred.JobClient:   File Output Format Counters 
INFO mapred.JobClient:     Bytes Written=249
INFO mapred.JobClient:   FileSystemCounters
INFO mapred.JobClient:     FILE_BYTES_READ=2253
INFO mapred.JobClient:     FILE_BYTES_WRITTEN=153667
INFO mapred.JobClient:   File Input Format Counters 
INFO mapred.JobClient:     Bytes Read=190
INFO mapred.JobClient:   Map-Reduce Framework
INFO mapred.JobClient:     Reduce input groups=35
INFO mapred.JobClient:     Map output materialized bytes=436
INFO mapred.JobClient:     Combine output records=0
INFO mapred.JobClient:     Map input records=16
INFO mapred.JobClient:     Reduce shuffle bytes=0
INFO mapred.JobClient:     Physical memory (bytes) snapshot=0
INFO mapred.JobClient:     Reduce output records=35
INFO mapred.JobClient:     Spilled Records=78
INFO mapred.JobClient:     Map output bytes=346
INFO mapred.JobClient:     CPU time spent (ms)=0
INFO mapred.JobClient:     Total committed heap usage (bytes)=728760320
INFO mapred.JobClient:     Virtual memory (bytes) snapshot=0
INFO mapred.JobClient:     Combine input records=0
INFO mapred.JobClient:     Map output records=39
INFO mapred.JobClient:     SPLIT_RAW_BYTES=216
INFO mapred.JobClient:     Reduce input records=39

Finalmente, este es el resultado del MapReduce:
Este 1
This 1
a 2
an 1
are 1
contar 1
count 1
cuantas 1
de 1
del 2
dentro 1
ejemplo 1
es 1
estan 1
example 1
fichero 2
going 1
hay 1
in 1
is 1
it 1
of 1
palabras 2
que 1
sus 1
text 1
that 1
the 1
to 1
un 1
vamos 1
ver 1
we 1
words 1
y 1

miércoles, mayo 07, 2014

Instalando Hadoop en una Raspberry Pi

¿Hadoop en una Raspberry Pi?

En cierto modo, se podría encontrar el motivo de esta entrada entre lo absurdo y lo innecesario. La capacidad de computación de una Raspberry Pi ni su memoria RAM no la hace excesivamente interesante para trabajos de Big Data, pero es cierto que ofrece una manera barata de montar un cluster con fines educativos en casa (especialmente si en casa tienes dos Raspberry Pi v2 una Raspberry Pi v1 y una BeagleBone Black)

La instalación de Hadoop en una Raspberry Pi resulta más sencillo de lo que cabe esperar. Y es que el pequeño ordenador de open hardware no se diferencia tanto de cualquier portátil u ordenador de sobremesa si obviamos la RAM de que dispone la v2 (512mb) y la capacidad del microprocesador (bueno, y el almacenamiento en sd, etc etc).

Básicamente, siguiendo el primero de los tutoriales sobre instalación de Hadoop e incluso modo tienes el 
90% del trabajo hecho. Sólo falta un pequeño detalle....

Y es que los trabajos de computación que se llegan a realizar en un Job MapReduce suelen requerir de gran cantidad de RAM y como nuestra pequeña Rpi no es que vaya sobrada precisamente, debemos indicar a la instalación de Hadoop unos límites para que el sistema operativo no se cargue los procesos del DataNode y del TaskTracker.

Esto es igualmente válido para ordenadores de consumo que no son los típicos que se pueden encontrar en un cluster de Hadoop de producción con chorrocientos Gb de RAM. En algún caso nos podríamos ver forzados a establecer límites en máquinas con más RAM, en alguna prueba que he hecho he tenido que establecer límites de 3gb en una máquina con 4Gb de RAM de 64 bits.

Déjate de rollos y dime que tengo que hacer

Fácil, para establecer estos límites tenemos que entrar en carpeta de configuración de Hadoop que en nuestra versión es $HADOOP_INSTALL/conf aunque también podría ser /etc/hadoop o en $HADOOP_INSTALL/etc/conf.

Una vez localizada la carpeta, abrir el archivo mapred-site.xml e introducir las siguientes propiedades:


    mapred.child.java.opt
    -Xmx386



    mapred.tasktracker.map.tasks.maximum
    1



    mapred.tasktracker.reduce.tasks.maximum
    1

Y voilá! Aunque parezca increíble esto fue lo único que tuve que hacer para agregar la Raspberry Pi a mi cluster casero. La explicación no es compleja: la primera propiedad limita la RAM utilizada por el proceso de Map o Reduce de la Rpi a 386, lo cual le da un pequeño margen para que el sistema operativo no nos mate el proceso. Las otras dos opciones limitan el número de operaciones de Map o Reduce que se pueden ejecutar en paralelo dentro de cada nodo de trabajo. Y es que si cada trabajo de Map o de Reduce es computacionalmente no muy intenso, se puede ajustar este parámetro para ejecutar varios trabajos en paralelo. En general he comprobado que en la Rpi no suele ser posible aumentar en número de tareas por encima de uno incluso en pequeños trabajos.

Y puestos a hablar de clusters caseros de Raspberrys Pi, ahí dejo un link de un señor que tuvo las pelotas.... Por la ley del máximo esfuerzo, de montarse un cluster con 40 Rpi Raspberry Pi 40 nodes cluster

lunes, mayo 05, 2014

Primeros pasos con Hive. Oootro contador de palabras

Descargando e instalando Hive

  1. Página web oficial: http://hive.apache.org/
  2. Site de releases: http://www.apache.org/dyn/closer.cgi/hive/
  3. Release que vamos a usar: http://ftp.cixug.es/apache/hive/hive-0.13.0/apache-hive-0.13.0-bin.tar.gz
Lo primero descargar y ejecutar la release que queramos usar:

wget 'http://ftp.cixug.es/apache/hive/hive-0.13.0/apache-hive-0.13.0-bin.tar.gz'
tar -zxvf apache-hive-0.13.0-bin.tar.gz
cd apache-hive-0.13.0-bin/bin
export HIVE_HOME=$HOME/apache-hive-q0.13.0-bin
export PATH=$PATH:$HIVE_HOME/bin

Arrancando el cluster de Hadoop

Hive necesita al menos Hadoop 1.2.1 corriendo para funcionar. Si no sabes como instalar Hadoop aún te recomiendo que te pases por el tutorial Instalando Hadoop en modo pseudo-distribuido (local)

Una vez tengamos el cluster corriendo, Hive necesita de dos carpetas dentro del HDFS con permisos de grupo para poder funcionar, para crear esas carpetas y darles los permisos necesarios ejecutaremos las órdenes siguientes:
hadoop fs -mkdir /user
hadoop fs -mkdir /user/hive
hadoop fs -mkdir /user/hive/warehouse
hadoop fs -mkdir /tmp
hadoop fs -chmod g+w /tmp
hadoop fs -chmod g+x /user/hive/warehouse

Como se puede apreciar, Hive necesita de las carpetas /tmp y /user/hive/warehouse para poder funcionar.

Ahora vamos a crear un fichero de ejemplo con el siguiente contenido para contar sus palabras. Lo llamaremos ejemplo.txt y lo vamos a guardar dentro de la carpeta bin de la carpeta de Hive:
Este es
un fichero de
ejemplo en
del que vamos
a contar
sus palabras
El cual tendremos que copiar dentro del HDFS. Recordamos el comando:
hadoop fs -copyFromLocal ejemplo.txt /tmp/ejemplo.txt
Ahora, ejecutamos la consola de Hive con el comando "hive" (sin comiillas) y creamos una tabla para almacenar el fichero de texto en ella:
hive>CREATE TABLE ejemplo (linea STRING);
Con el comando anterior, sólo hemos creado una tabla nueva, pero no la hemos cargado de información. Para cargar un fichero en una tabla ejecutaremos la orden siguiente:
hive>LOAD DATA INPATH '/tmp/ejemplo.txt' OVERWRITE INTO TABLE ejemplo;

Deleted hdfs://asusnotebook:54310/user/hive/warehouse/texto

Table default.texto stats: [numFiles=1, numRows=0, totalSize=73, rawDataSize=0]

OK

Time taken: 1.612 seconds
Ya tenemos el fichero cargado en la tabla ejemplo. Ahora hay que crear otra tabla para almacenar el resultado del conteo de palabras:

hive>CREATE TABLE contador AS 
SELECT palabra, count(1) AS cuenta FROM 
(SELECT explode(split(linea,' ')) AS palabra FROM texto) w 
GROUP BY palabra 
ORDER BY cuenta;
Esto comenzará creará dos MapReduce que comenzarán a trabajar de inmediato. Los mensajes que aparecen son como los siguientes (resumido)
[...]

Stage-2 map = 0%,  reduce = 0%

Stage-2 map = 100%,  reduce = 0%, Cumulative CPU 0.84 sec

Stage-2 map = 100%,  reduce = 100%, Cumulative CPU 2.12 sec

OK

Time taken: 42.304 seconds

Vamos a ver ahora que se ha creado dentro de la tabla contador

hive>SELECT * FROM contador;

a 1
contar 1
cuantas 1
de 1
del 2
dentro 1
ejemplo 1
es 1
estan 1
fichero 2
hay 1
palabras 2
que 1
sus 1
un 1
vamos 1
ver 1
y 1
Y ahí lo tenemos. Por poco casi no repito una palabra pero como se puede apreciar, la palabra "del", "fichero" y "palabras" aparecen 2 veces cada una.

lunes, abril 28, 2014

Escribiendo scripts en Apache Pig para Hadoop

Descargando y ejecutando Pig en modo local

Si aún no has hecho el primero de los tutoriales sobre Pig, te recomiendo pasarte por él ya que va a ser el mínimo necesario para que puedas realizar este tutorial: Primeros pasos con Apache Pig. Usando Pig para hacer un contador de palabras

Descripción breve de un script Pig

Un script Pig no es mas que una serie de órdenes ejecutadas secuencialmente para realizar una consulta. Es exactamente lo mismo que utilizar la consola sólo que además le puedes poner comentarios. La manera de escribir comentarios en el código Pig es precediendo las líneas con un doble guión (--) o para hacerlo multilínea sería con el típico /* */

El contador de palabras en un script Pig

Para hacer el contador de palabras en un script Pig, lo único que tenemos que hacer es crear un fichero con la extensión pig e introducir las órdenes que utilizamos en el tutorial anterior sobre Pig en el fichero. Vamos a crear un fichero llamado wordcounter.pig y dentro escribiremos lo siguiente:

/* Cargar el fichero */
myfile = LOAD 'tweets' AS (words:chararray);

/* Separar las palabras dentro de cada linea */
wordsList = FOREACH myfile GENERATE TOKENIZE($0);

/* Separar las palabras a una linea por palabra */
words = FOREACH wordsList GENERATE FLATTEN($0);       

/* Agrupar las palabras iguales en la misma linea */
groupedWords = GROUP words BY $0;

/* Contar las palabras */
final = FOREACH groupedWords GENERATE $0, COUNT($1);  

/* Ordenar las palabras */
sortedWords = ORDER final BY $1 ASC;

/* Mostrar los resultados por pantalla */
DUMP sortedWords;


Ahora, para ejecutar el script en modo local, sólo tenemos que escribir:
pig -x local wordcounter.pig

Bonus: Guardando los resultados en un archivo

En la mayoría de los casos no vamos a querer mostar los resultados por pantalla sino que vamos a querer que los resultados se guarden en disco para su posterior análisis. Nada mas fácil en Pig. Lo único que tenemos que hacer es reemplazar la instrucción DUMP con la instrucción STORE. El código nos quedaría así:


/* Cargar el fichero */
myfile = LOAD 'tweets' AS (words:chararray);

/* Separar las palabras dentro de cada linea */
wordsList = FOREACH myfile GENERATE TOKENIZE($0);

/* Separar las palabras a una linea por palabra */
words = FOREACH wordsList GENERATE FLATTEN($0);       

/* Agrupar las palabras iguales en la misma linea */
groupedWords = GROUP words BY $0;

/* Contar las palabras */
final = FOREACH groupedWords GENERATE $0, COUNT($1);  

/* Ordenar las palabras */
sortedWords = ORDER final BY $1 ASC;

/* Guardar los resultados en una carpeta llamada pig_wordcount */
STORE sortedWords into 'pig_wordcount';


Esto creará una carpeta llamada "pig_wordcount" que contendrá los resultados de la consulta en archivos con un formato "part-r-*". Para ver los resultados sólo habrá que ejecutar desde el bash:
cat pig_wordcount/part-r-*

sábado, abril 26, 2014

Primeros pasos con Apache Pig. Usando Apache Pig para hacer un Wordcount en modo local

Descargando e instalando Pig

Algunos links de interés:

Lo primero será descargar, extraer Pig y decirle donde está la instalación de Hadoop (ver tutorial Instalando Hadoop en modo pseudo-distribuido (esto no es estrictamente necesario en el modo local pero nos ahorrará tiempo y quebraderos de cabeza en el futuro):
 
wget 'http://apache.rediris.es/pig/pig-0.12.1/pig-0.12.1.tar.gz'
tar -zxvf pig-0.12.1.tar.gz
cd pig-0.12.1/bin
export PIG_CLASSPATH=$HADOOP_INSTALL


Bien, ahora vamos a hacer el mismo contador de palabras del tutorial anterior pero con Pig, para eso vamos a descargar el mismo archivo en formato de texto plano http://www.gutenberg.org/ebooks/11101 y lo vamos a poner en la carpeta bin dentro de Pig. Deberíamos tener los siguientes ficheros dentro de la carpeta de Pig:
total 204K
-rwxr-xr-x 1 fedora fedora 13K abr 5 10:43 pig
-rw-r--r-- 1 fedora fedora 161K abr 19 23:51 pg11101.txt
-rwxr-xr-x 1 fedora fedora 5,6K abr 5 10:43 pig.cmd
-rwxr-xr-x 1 fedora fedora 14K abr 5 10:43 pig.py
/home/fedora/Descargas/pig-0.12.1/bin
[fedora@localhost bin]$


Ejecutando Pig en modo local

pig -x local
Despues de unos cuantos mensajes de log aparece grunt
grunt>

Ahora vamos a escribir el contador de palabras en Pig, lo primero será cargar el archivo con la instrucción LOAD:
grunt>myfile = LOAD 'pg11101.txt' AS (words:chararray);

Si usamos la instrucción DUMP para ver los contenidos de la variable myfile, podemos ver la estructura general que tiene ahora cada linea:
grunt>DUMP myfile;
...
(EBooks posted since November 2003, with etext numbers OVER #10000, are)
(filed in a different way. The year of a release date is no longer part)
(of the directory path. The path is based on the etext number (which is)
...

Todavía falta bastante, ni siquiera tenemos una lista de palabras. Para poder separarlo por palabras, podemos usar la instrucción TOKENIZE:
grunt> wordsList = FOREACH myfile GENERATE TOKENIZE($0);
...
({(EBooks),(posted),(since),(November),(2003),(with),(etext),(numbers),(OVER),(#10000),(are)})
({(filed),(in),(a),(different),(way.),(The),(year),(of),(a),(release),(date),(is),(no),(longer),(part)})
({(of),(the),(directory),(path.),(The),(path),(is),(based),(on),(the),(etext),(number),(which),(is)})
...

Si nos fijamos bien, ahora tenemos tuplas (arrays) de palabras pero seguimos sin tener un listado de palabras separadas: Para conseguir dicho listado, tenemos la instrucción FLATTEN:
grunt>words = FOREACH wordsList GENERATE FLATTEN($0);
...
(at:)
()
(http://www.gutenberg.net/1/0/2/3/10234)
()
(or)
(filename)
(24689)
(would)
(be)como
(found)
(at:)
...

Con lo que ahora tenemos las listas separadas, una palabra por linea (aunque todavía están repetidas). Ahora es el momento de agruparlas usando la instrucción GROUP:
grunt>groupedWords = GROUP words BY $0;
...
(www.gutenberg.net,{(www.gutenberg.net),(www.gutenberg.net),(www.gutenberg.net)})
(Constantinopolitan,{(Constantinopolitan),(Constantinopolitan)})
...

La cosa marcha, ahora tenemos listados agrupados de la misma palabra por cada línea. Sólo nos falta contar cuantas ocurrencias hay por línea. Para ello le diremos que por cada linea nos tiene que crear una nueva pareja clave/valor que clave ($0) siga siendo el nombre de la palabra que estamos contando, como hasta ahora y el valor la cuenta de ítems en la segunda posición ($1) con la instrucción COUNT.
grunt>final = FOREACH groupedWords GENERATE $0, COUNT($1);
...
(www.gutenberg.net,3)
(Constantinopolitan,2)
...

Y ahí lo tenemos, nuestro conteo de palabras.


Bonus: Ordenando ocurrencias en Pig

Y como extra, vamos a ordenar el resultado del conteo de palabras en orden descendente, de tal manera que el último de los resultados de la lista sea la palabra mas usada de todo el texto:
grunt>sortedWords = ORDER final BY $1 ASC;
...
(a,697)
(in,833)
(and,1137)
(of,1323)
(the,2047)
...

Por fin, parece que nuestro ganador es la palabra "the". Espero que el tutorial haya resulado interesante :)

miércoles, abril 23, 2014

El Ecosistema de Hadoop - Un resumen de las tecnologías que envuelven Hadoop

El ecosistema de Hadoop

O como tener un montón de tecnologías juntas con un objetivo común

Cuando empecé a aprender Hadoop, una de las primeras cosas que me pasaron fue que me sentía completamente sobrepasado por el nivel de información y términos que rodean Hadoop. HDFS, MapReduce, Streaming, Pig, Hive... puede que incluso frases como "commodity hardware" puedan resultar desconocidas.

Y es que Hadoop es un completo ecosistema de soluciones con un objetivo común, la explotación de grandes cantidades de datos.

Así que, para intentar aclarar un poco las cosas, comparto un link en el que se puede leer un resumen de todo el ecosistema de Hadoop (eso si, en inglés, intentaré traducirlo al español si saco un hueco).

¡Saludos!

lunes, abril 21, 2014

Big-Data con Hadoop

La ley del máximo esfuerzo


A lo largo de las próximas semanas, voy a estar haciendo un poco de investigación a bajo nivel de Hadoop y todo su ecosistema. Algunos se preguntarán (y con razón), ¿porqué motivo reinventar la rueda cuando hay soluciones excelentes ya montadas para la explotación del Big Data como los que ofrecen Cloudera (http://www.cloudera.com/), Horton Works (hortonworks.com) o IBM InfoSphere BigInsights (http://www-01.ibm.com/software/data/infosphere/biginsights)?

Bueno, por la ley del máximo esfuerzo, por supuesto, aquella por la cual cuando tienes sed caminas decenas (o cientos) de kilómetros para ir al mar, llenar un vaso con agua, desalinizarlo y bebértelo... sólo para darte cuenta de que sigues teniendo sed... pero ya sabes como funciona un proceso de desalinización.

NOTA: Al hilo de la ley del máximo esfuerzo, hay un excelente libro de un británico que se hizo una tostadora de cero a base de conseguir los materiales que necesitaba, trabajarlos y montarlos. El libro en cuestión se llama "The Toaster Project" de Thomas Thwaites y se puede adquirir en formato kindle en The Toaster Project en Amazon.es

Tecnologías Big Data

La investigación se va a centrar, en este orden mas o menos, en los siguientes puntos:

  1. Instalando Apache Hadoop en modo Pseudo-Distribuido
  2. El Ecosistema de Hadoop - Un resumen de las tecnologías que envuelven Hadoop
  3. Escribiendo un Hadoop MapReduce en Java: Un Wordcount mejorado
  4. Instalando Apache Hadoop en modo Distribuido
  5. Usando Amazon EMR para ejecutar un Hadoop MapReduce
  6. Primeros pasos con Apache Pig. Usando Pig para hacer un contador de palabras
  7. Usando una Rapsberry Pi como esclavo en Apache Hadoop
  8. Creando scripts en Apache Pig para Hadoop
  9. Usando Apache Pig en modo distribuido (o pseudo-distribuido)
  10. Usando Amazon EMR para ejecutar nuestras consultas con Apache Pig
  11. Primeros pasos con Apache Hive. Creando el contador de palabras en Apache Hive
  12. Aplicaciones distribuidas con Zoekeeper.
  13. Real Time Analytics con Apache Storm: Un contador de palabras en tiempo real
  14. Instalando un cluster de Apache Storm
  15. Creando scripts de Apache Hive y ejecutándolos en el cluster de Hadoop
  16. Usando Amazon EMR para ejecutar nuestras consultas en Apache Hive
  17. Primeros pasos con machine-learning usando Mahout... un recomendador de productos
  18. Creando flujos de trabajos con Oozie
  19. Introducción a la logística de datos con Apache Flume
  20. Logística de datos con Apache Sqoop: Importando bases de datos al HDFS
  21. Logística de datos con Apache Sqoop: Exportando datos del HDFS a MySQL
  22. Analizando datos con R y Hadoop.
  23. Introducción a Apache Spark
  24. Introducción y primeros pasos con Hbase
  25. Introducción a Apache Kafka: un recolector de mensajes
  26. Tanteando Shark y sus posibilidades
  27. Mas machine learning con MLlib
  28. Escapando del modo consola con GraphX. Gráficos para Big Data
  29. Haciendo dashboards de Big Data con Intellicus
  30. Mas...?

Páginas webs oficiales

Apache Hadoop: hadoop.apache.org/‎
Amazon AWS: aws.amazon.com/‎
Apache Hive: hive.apache.org
Apache Flume: flume.apache.org
Apache Mahout: http://mahout.apache.org

sábado, abril 19, 2014

Instalando Hadoop en modo pseudo-distribuido

Bueno, lo primero de todo es que yo no soy ningún experto en Hadoop por lo que el post va más orientado a mis "notas personales" sobre cómo instalé un sistema pseudo-distribuido de Hadoop. En el siguiente post explicaré un método para ejecutar Hadoop en tu propia máquina "engañando" (aunque de engañar nada) a la máquina para que piense que tiene un nodo "esclavo".

Primero un poco de terminología básica y MUY resumida en modo "Entendimiento de abuela":
  1. Nodo "Maestro":
    1. Namenode: Nodo jefe si se le quiere llamar así. Es el que maneja el meollo de todo.
    2. Jobtracker: Se encarga de asignar las tareas de computación.
    3. SecondaryNamenode: Tantea que todo esté funcionando correctamente de vez en cuando.
  2. Nodo "Esclavo":
    1. Tasktracker: Busca tareas de computación para realizar.
    2. Datanode: Es el proceso que hace el trabajo "sucio". Osea, el que "computa" realmente.
Y básicamente, el sistema pseudo-distribuido lo que hace es tener todos esos procesos en la misma máquina.

Para este tutorial vamos a usar una versión que ha estado siendo usada en producción muy a menudo. La versión 1.2.1.

Pre-requisitos

  1. Java 6
  2. Yum
  3. SSH sin contraseña
Mi máquina es una Fedora por lo que en general usaré la terminología de Fedora a la hora de instalar paquetes. Afortunadamente, para trasladar esto a Ubuntu o similares basta con cambiar la palabra "yum" con la palabra "apt-get". Para instalar Java simplemente hay que ejecutar:
sudo yum update
sudo yum install -y opendjk-6-jdk
echo "export JAVA_HOME=/usr/lib/jvm/openjdk-6-jdk" >> ~/.bashrc
Osea, actualizar el sistema. Instalar Java 6 y exportar la variable de entorno JAVA_HOME en caso de que no se hubiera exportado automáticamente.

Hay que habilitar el acceso por SSH sin contraseña a los nodos esclavos. En este caso el nodo esclavo es la misma máquina así que tenemos que permitir un acceso por SSH sin contraseña... ¡a nosotros mismos!

ssh-keygen -t rsa -P ""

Preguntará la ubicación del archivo que, generalmente, irá en ~/.ssh/id_rsa. Mi consejo es llamarlo id_rsa_passless para diferenciarlo del resto de claves rsa que podamos tener. Ahora habrá que agregarla a nuestra lista de claves autorizadas:
cat ~/.ssh/id_rsa_passless.pub >> ~/.ssh/authorized_keys

Ahora, si hacemos ssh localhost, si es la primera vez nos pregutará si confiamos en la máquina, le decimos que "yes" y ya podremos entrar si contraseña.

Descargando Hadoop 

Descargar la version 1.2.1 desde la web oficial http://hadoop.apache.org/. NOTA: No hace falta descargar la versión src, con la normal vale. Extraer el contenido del archivo comprimido
tar -zxvf hadoop-1.2.1.tar.gz
 Mover la carpeta extraída a la localización final de Hadoop. Ojo, la localización de la carpeta es importante ya que luego tendremos que hacer referencia a ella.
sudo mv hadoop-1.2.1.tar.gz   /var/hadoop
Asignar el usuario y grupo correspondiente:
sudo chown -R [user] /var/hadoop
sudo chgrp -R [user] /var/hadoop
mkdir /var/hadoop/hdfs
echo "export HADOOP_INSTALL=/var/hadoop" >> ~/.bashrc
echo "export HADOOP_OPTS=-Djava.net.preferIPv4Stack=true" >> ~/.bashrc
Donde [user] es el nombre del usuario que esté ejecutándose en la máquina. Algunos tutoriales aconsejan hacer un usuario hduser y un grupo hadoop. Esto es recomendable pero para hacer las cosas lo mas simples posibles he preferido omitirlo. Por último creamos la carpeta hdfs dentro del directorio de Hadoop para usarla como el sistema de ficheros HDFS (Hadoop Distributed File System). Por último exportamos el directorio de instalación de Hadoop y le indicamos que sólo use direcciones IPv4.

Configurar Hadoop

Hadoop consta, en la versión de producción mas usada, de tres ficheros de configuraciones fundamentales. Estos los vamos a tener de esta forma:

core-site.xml: 

Configurar la carpeta para el sistema de ficheros HDFS y configurar la máquina del Namenode.

<?xml version="1.0"?>
<?xml-stylesheet type="text/xsl" href="configuration.xsl"?>
<configuration>
<property>
  <name>hadoop.tmp.dir</name>
  <value>/var/hadoop/hdfs</value>
</property>

<property>
  <name>fs.default.name</name>
  <value>hdfs://localhost:54310</value>
</property>
</configuration>

hdfs-site.xml

Para este fichero le estamos diciendo que:

  1. Use sólo una replicación por bloque (por ejemplo, con 3 replicaciones cada bloque de datos existe 3 veces en el cluster para, en caso de fallo, tener una copia del mismo).
  2. No use permisos para el sistema de ficheros HDFS. Esto NO debería hacerse así, pero por simplicidad lo dejamos así que nos va a ahorrar algunos problemas.
  3. No usar mas de 3072 mb de RAM. Esto es para evitar que el propio sistema linux mate el proceso si ve que está bloqueando el resto del sistema. Para una máquina de 4Gb de RAM es un umbral razonable.
  4. Permitir un número de procesos hasta 4096. Esto evita ciertos errores que he tenido comúnmente que, al parecer, estaban relacionados con este valor

<?xml version="1.0"?>
<?xml-stylesheet type="text/xsl" href="configuration.xsl"?>
<configuration>

<property>
  <name>dfs.replication</name>
  <value>1</value>
</property>

<property>
    <name>dfs.permissions</name>
    <value>false</value>
</property>

<property>
<name>mapred.child.java.opts</name>
<value>-Xmx3072m</value>
</property>

<property>
<name>dfs.datanode.max.xcievers</name>
<value>4096</value>
</property>

</configuration>

mapred-site.xml

En el fichero Mapred le vamos a indicar en qué máquina está el JobTracker:

<?xml version="1.0"?>
<?xml-stylesheet type="text/xsl" href="configuration.xsl"?>
<configuration>
<property>
  <name>mapred.job.tracker</name>
  <value>localhost:54311</value>
</property>
</configuration>

masters

Contendrá un listado de máquinas para usar el SecondaryNamenode en ellas:
localhost

slaves

Contendrá un listado de máquinas que van a efectuar trabajos de computación. Como el nuestro es un sistema pseudo-distribuido, es el mismo que el namenode, secondarynamenode, etc:
localhost 

Formatear el sistema de ficheros HDFS

Aunque tenga la palabra formatear y asuste mucho, en realidad lo único que formatea son los contenidos de la carpeta que hayamos definido anteriormente en el fichero core-site.xml, la propiedad hadoop.tmp.dir. En nuestro caso es la carpeta /var/hadoop/hdfs
cd /var/hadoop
bin/hadoop namenode -format

 Copiar un archivo para ejecutar el trabajo sobre él

Para nuestro ejemplo, vamos a usar el poema de Samuel Taylor Coleridge "Ancient Mariner and selected poems" disponible en el Proyecto Gutenberg -> http://www.gutenberg.org/ebooks/11101 (descargar la versión plain text). Movemos el archivo descargado a la carpeta /var/hadoop y creamos una carpeta en el HDFS para copiar el archivo dentro.
bin/hadoop fs -mkdir /rime
bin/hadoop fs -copyFromLocal pg11101.txt /rime

Iniciar el pseudo-cluster

Ahora, para iniciar el cluster y que toda la magia se ponga en marcha hay que ejecutar el siguiente comando:
sbin/start-all.sh
Podemos comprobar si todo ha ido bien escribiendo en el terminal jps lo cual debe darnos un resultado similar al siguiente:
6246 Jps
2435 DataNode
2740 JobTracker
2313 NameNode
2605 SecondaryNameNode
2875 TaskTracker
Si algo ha ido mal, se puede comprobar el log que se almacenará en la carpeta logs del directorio de hadoop.

Todo listo: Iniciar el trabajo del contador de palabras

Todos los procesos corriendo, el fichero de prueba en el HDFS....
bin/hadoop jar share/hadoop/mapreduce/hadoop-mapreduce-examples-0.23.9.jar wordcount /rime/pg11101.txt /rime/output
Osea:

  1. bin/hadoop = Ejecuta el binario de Hadoop...
  2. jar = ...para ejecutar el jar... 
  3. share/hadoop/mapreduce/hadoop-mapreduce-examples-0.23.9.jar = (jar a ejecutar) del cual usa el subproceso...
  4. wordcount = ...subproceso a ejecutar...
  5. /rime/pg11101.txt = ...para contar este fichero de mi HDFS...
  6. /rime/output = ...y guardarme el resultado en esta carpeta.
Este es un resultado típico de ejecutar un mapreduce:

INFO input.FileInputFormat: Total input paths to process : 1
INFO mapred.JobClient: Running job: job_201404191919_0002
INFO mapred.JobClient: map 0% reduce 0%
INFO mapred.JobClient: map 100% reduce 0%
INFO mapred.JobClient: map 100% reduce 100%
INFO mapred.JobClient: Job complete: job_201404191919_0002
INFO mapred.JobClient: Counters: 29
INFO mapred.JobClient: Job Counters
INFO mapred.JobClient: Launched reduce tasks=1
INFO mapred.JobClient: SLOTS_MILLIS_MAPS=13320
INFO mapred.JobClient: Total time spent by all reduces waiting after reserving slots (ms)=0
INFO mapred.JobClient: Total time spent by all maps waiting after reserving slots (ms)=0
INFO mapred.JobClient: Launched map tasks=1
INFO mapred.JobClient: Data-local map tasks=1
INFO mapred.JobClient: SLOTS_MILLIS_REDUCES=10400
INFO mapred.JobClient: File Output Format Counters
INFO mapred.JobClient: Bytes Written=103002
INFO mapred.JobClient: FileSystemCounters
INFO mapred.JobClient: FILE_BYTES_READ=143849
INFO mapred.JobClient: HDFS_BYTES_READ=231867
INFO mapred.JobClient: FILE_BYTES_WRITTEN=330699
INFO mapred.JobClient: HDFS_BYTES_WRITTEN=103002
INFO mapred.JobClient: File Input Format Counters
INFO mapred.JobClient: Bytes Read=231760
INFO mapred.JobClient: Map-Reduce Framework
INFO mapred.JobClient: Map output materialized bytes=143849
INFO mapred.JobClient: Map input records=4989
INFO mapred.JobClient: Reduce shuffle bytes=143849
INFO mapred.JobClient: Spilled Records=20624
INFO mapred.JobClient: Map output bytes=356537
INFO mapred.JobClient: Total committed heap usage (bytes)=251133952
INFO mapred.JobClient: CPU time spent (ms)=3950
INFO mapred.JobClient: Combine input records=36475
INFO mapred.JobClient: SPLIT_RAW_BYTES=107
INFO mapred.JobClient: Reduce input records=10312
INFO mapred.JobClient: Reduce input groups=10312
INFO mapred.JobClient: Combine output records=10312
INFO mapred.JobClient: Physical memory (bytes) snapshot=267968512
INFO mapred.JobClient: Reduce output records=10312
INFO mapred.JobClient: Virtual memory (bytes) snapshot=10779578368
INFO mapred.JobClient: Map output records=36475

Y ya está

Si todo ha ido bien el siguiente comando nos mostrará el conteo de palabras del fichero que le hemos pasado:
bin/hadoop fs -cat /rime/output/part-r-00000
wrath 1
wrath*, 1
wreathless 1
wreaths 1
wrenched 1
wretch! 1
wretched 2
wretchedness, 1
write 3
writes 3
writing 6
writings 3
written 17
written, 2
written. 1
wrong 2
wrong, 1
wrong; 1
wronged 1
wrote 25
wrote, 1
wrote: 2
wroth 1
www.gutenberg.net 2
xi; 1
xiv.). 1
xvii., 1
ye 10
year 13
year, 9
year. 3
year: 1
year; 1
Como se puede apreciar en el resultado de ejemplo, quedaría una última vuelta de tuerca en el programa MapReduce del Wordcount en el cual habría que intentar agrupar las palabras con contuvieran únicamente caracteres alfanuméricos, esto lo haremos en el tutorial sobre como escribir nuestro propio MapReduce: Un Wordcount mejorado.