Ejemplo de Hadoop MapReduce: Primero Java Programa con Code

⚡ Resumen inteligente

Los programas Hadoop MapReduce se escriben como tres Java clases, un mapeador, un reductor y un controlador, que se compilan, se empaquetan en un archivo jar y se envían al clúster para contar las ventas por país.

  • 🔘 Conjunto de datos: El archivo SalesJan2009.csv contiene una transacción por línea, con el país en la octava columna separada por comas.
  • ☑️ Mapeador: SalesMapper divide cada línea y emite el país emparejado con el valor constante uno.
  • Reductor: SalesCountryReducer suma los que llegan para cada país en un único recuento de frecuencia.
  • 🧪 Conductor: SalesCountryDriver nombra el trabajo, declara los tipos de clave y valor y conecta ambas clases entre sí.
  • 🛠️ Build: Compila con javac, agrega una entrada de manifiesto Main-Class y luego empaqueta todo con jar cfm.
  • ⚠️ Ejecutar: Copia el archivo CSV en HDFS, envía el archivo jar y lee la parte 00000 del directorio de salida.

Ejemplo de Hadoop MapReduce creando un primero Java programa

En este tutorial, aprenderá a utilizar Hadoop con ejemplos de MapReduce. Los datos de entrada utilizados son VentasEne2009.csvContiene información relacionada con las ventas, como el nombre del producto, el precio, la forma de pago, la ciudad y el país del cliente. El objetivo es determinar la cantidad de productos vendidos en cada país.

Primer programa Hadoop MapReduce

Ahora en esto Tutorial de MapReduce, crearemos nuestro primer Java Programa MapReduce:

La captura de pantalla que aparece a continuación muestra los datos brutos de Ventas de enero de 2009, donde cada línea representa una transacción y el país se encuentra en la octava columna separada por comas.

Los datos de ventas de enero de 2009 se abrieron en una hoja de cálculo que mostraba las columnas de transacciones.

Asegúrese de tener Hadoop instalado. Antes de comenzar con el proceso, cambie el usuario a 'hduser' (el ID utilizado durante la configuración de Hadoop; puede cambiar al ID de usuario utilizado durante su propia configuración de Hadoop).

su - hduser_

El mensaje cambia a la cuenta hduser, como se muestra a continuación.

Símbolo del terminal después de cambiar a la cuenta hduser

Paso 1) Crea el directorio del proyecto y los archivos fuente.

Cree un nuevo directorio con el nombre MapReduceTutorial, como se muestra en el siguiente ejemplo de MapReduce.

sudo mkdir MapReduceTutorial

El comando mkdir crea el directorio MapReduceTutorial.

Dar permisos

sudo chmod -R 777 MapReduceTutorial

Comando chmod que otorga permisos completos en MapReduceTutorial

Crea los tres Java Los archivos fuente se encuentran a continuación dentro de MapReduceTutorial. Tenga en cuenta que los tres utilizan la versión anterior. org.apache.hadoop.mapred API, que todavía viene incluida con Hadoop 3.x.

SalesMapper.java

package SalesCountry;

import java.io.IOException;

import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapred.*;

public class SalesMapper extends MapReduceBase implements Mapper <LongWritable, Text, Text, IntWritable> {
	private final static IntWritable one = new IntWritable(1);

	public void map(LongWritable key, Text value, OutputCollector <Text, IntWritable> output, Reporter reporter) throws IOException {

		String valueString = value.toString();
		String[] SingleCountryData = valueString.split(",");
		output.collect(new Text(SingleCountryData[7]), one);
	}
}

SalesCountryReducer.java

package SalesCountry;

import java.io.IOException;
import java.util.*;

import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapred.*;

public class SalesCountryReducer extends MapReduceBase implements Reducer<Text, IntWritable, Text, IntWritable> {

	public void reduce(Text t_key, Iterator<IntWritable> values, OutputCollector<Text,IntWritable> output, Reporter reporter) throws IOException {
		Text key = t_key;
		int frequencyForCountry = 0;
		while (values.hasNext()) {
			// replace type of value with the actual type of our value
			IntWritable value = (IntWritable) values.next();
			frequencyForCountry += value.get();
			
		}
		output.collect(key, new IntWritable(frequencyForCountry));
	}
}

SalesCountryDriver.java

package SalesCountry;

import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.*;
import org.apache.hadoop.mapred.*;

public class SalesCountryDriver {
    public static void main(String[] args) {
        JobClient my_client = new JobClient();
        // Create a configuration object for the job
        JobConf job_conf = new JobConf(SalesCountryDriver.class);

        // Set a name of the Job
        job_conf.setJobName("SalePerCountry");

        // Specify data type of output key and value
        job_conf.setOutputKeyClass(Text.class);
        job_conf.setOutputValueClass(IntWritable.class);

        // Specify names of Mapper and Reducer Class
        job_conf.setMapperClass(SalesCountry.SalesMapper.class);
        job_conf.setReducerClass(SalesCountry.SalesCountryReducer.class);

        // Specify formats of the data type of Input and output
        job_conf.setInputFormat(TextInputFormat.class);
        job_conf.setOutputFormat(TextOutputFormat.class);

        // Set input and output directories using command line arguments, 
        //arg[0] = name of input directory on HDFS, and arg[1] =  name of output directory to be created to store the output file.

        FileInputFormat.setInputPaths(job_conf, new Path(args[0]));
        FileOutputFormat.setOutputPath(job_conf, new Path(args[1]));

        my_client.setConf(job_conf);
        try {
            // Run the job 
            JobClient.runJob(job_conf);
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

Descargar archivos aquí

El archivo se descompone en los mismos tres archivos fuente, como se muestra aquí.

ExtracListado de archivos de ted: SalesMapper, SalesCountryReducer y SalesCountryDriver

Verifique los permisos de archivo de todos estos archivos.

Listado extenso que muestra los permisos de archivo en los tres Java archivos fuente

Si faltan los permisos de "lectura", concédalos:

El comando chmod agrega permisos de lectura al Java archivos fuente

Paso 2) Exportar la ruta de clases de Hadoop

Exporta la ruta de clases como se muestra en el siguiente ejemplo de Hadoop.

export CLASSPATH="$HADOOP_HOME/share/hadoop/mapreduce/hadoop-mapreduce-client-core-2.2.0.jar:$HADOOP_HOME/share/hadoop/mapreduce/hadoop-mapreduce-client-common-2.2.0.jar:$HADOOP_HOME/share/hadoop/common/hadoop-common-2.2.0.jar:~/MapReduceTutorial/SalesCountry/*:$HADOOP_HOME/lib/*"

La ruta de clases exportada se muestra de nuevo en la línea de comandos.

Mensaje de la consola después de exportar la variable CLASSPATH de Hadoop.

Paso 3) Compile el Java archivos

Compila el Java archivos (estos archivos se encuentran en el directorio Final-MapReduceHandsOn). Sus archivos de clase se colocarán en el directorio del paquete.

javac -d . SalesMapper.java SalesCountryReducer.java SalesCountryDriver.java

Salida de javac que muestra una advertencia de obsolescencia durante la compilación.

Esta advertencia puede ignorarse sin problema; simplemente informa de que la API de mapred está obsoleta.

Esta compilación creará un directorio en el directorio actual con el nombre del paquete especificado en el Java archivo fuente (es decir SalesCountry en nuestro caso) y coloque todos los archivos de clase compilados en él.

Directorio del paquete SalesCountry que contiene los tres archivos de clase compilados.

Paso 4) Crear el archivo de manifiesto

Crea un nuevo archivo llamado Manifest.txt.

sudo gedit Manifest.txt

Añádele la siguiente línea:

Main-Class: SalesCountry.SalesCountryDriver

gedit Ventana que muestra la entrada Main-Class en Manifest.txt

SalesCountry.SalesCountryDriver es el nombre de la clase principal. Tenga en cuenta que debe presionar la tecla Enter al final de esta línea.

Paso 5) Empaqueta las clases en un frasco

Crear un archivo Jar

jar cfm ProductSalePerCountry.jar Manifest.txt SalesCountry/*.class

El comando jar crea ProductSalePerCountry.jar a partir del manifiesto.

Compruebe que el archivo jar esté creado

Se creó un listado en el directorio que confirma la creación de ProductSalePerCountry.jar.

Paso 6) Iniciar Hadoop

Iniciar Hadoop

$HADOOP_HOME/sbin/start-dfs.sh
$HADOOP_HOME/sbin/start-yarn.sh

Paso 7) Copie el archivo de entrada en HDFS.

Copia el archivo SalesJan2009.csv en ~/inputMapReduce

Ahora utilice el siguiente comando para copiar ~/inputMapReduce a HDFS.

$HADOOP_HOME/bin/hdfs dfs -copyFromLocal ~/inputMapReduce /

Salida del comando copyFromLocal con una advertencia de biblioteca nativa

Podemos ignorar con seguridad esta advertencia.

Verifique si un archivo realmente se copia o no.

$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce

hdfs dfs -ls listando SalesJan2009.csv dentro del directorio inputMapReduce

Paso 8) Ejecutar el trabajo MapReduce

Ejecute el trabajo MapReduce

$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales

Salida de la consola mientras se ejecuta el trabajo SalePerCountry en el clúster.

Esto creará un directorio de salida llamado mapreduce_output_sales en HDFS. El contenido de este directorio será un archivo que contendrá las ventas de productos por país.

Paso 9) Lea los resultados

El resultado se puede ver a través de la interfaz de comandos como,

$HADOOP_HOME/bin/hdfs dfs -cat /mapreduce_output_sales/part-00000

Pares de país y recuento de ventas impresos por el comando hdfs dfs -cat

Los resultados también se pueden ver a través de una interfaz web como:

Abra http://localhost:50070/ en un navegador web. En Hadoop 3.x, la interfaz web de NameNode se movió al puerto 9870, así que usa http://localhost:9870/ en lugar de eso

Página de inicio de la interfaz web de NameNode en un navegador

Ahora seleccione 'Explorar el sistema de archivos' y navegue hasta /mapreduce_output_sales

Explorador de archivos HDFS que muestra el directorio mapreduce_output_sales

Abrir parte-r-00000

El archivo de resultados part-r-00000 se abrió en el navegador HDFS.

Explicación de la clase SalesMapper

A medida que el proceso se desarrolla de principio a fin, las siguientes tres secciones explican en detalle qué hace cada clase.

En esta sección, comprenderemos la implementación de la clase SalesMapper.

1. Comenzamos especificando un nombre de paquete para nuestra clase. SalesCountry es el nombre de nuestro paquete. Tenga en cuenta que el resultado de la compilación, SalesMapper.class, se guardará en un directorio con el mismo nombre de paquete: SalesCountry.

Seguido de esto, importamos paquetes de biblioteca.

La siguiente captura de pantalla muestra una implementación de la clase SalesMapper.

Vista del editor de la implementación completa de la clase SalesMapper

Muestra Code Explicación:

1. Definición de clase SalesMapper

public class SalesMapper extends MapReduceBase implements Mapper<LongWritable, Text, Text, IntWritable> {

Cada clase de mapeador debe heredar de la clase MapReduceBase y debe implementar la interfaz Mapper.

2. Definición de la función 'mapa'

public void map(LongWritable key,
         Text value,
OutputCollector<Text, IntWritable> output,
Reporter reporter) throws IOException

La parte principal de la clase Mapper es un método 'map()' que acepta cuatro argumentos.

En cada llamada al método 'map()', se pasa un par clave-valor ('clave' y 'valor' en este código).

El método 'map()' comienza dividiendo el texto de entrada que se recibe como argumento. Divide cada línea en campos.

String valueString = value.toString();
String[] SingleCountryData = valueString.split(",");

Aquí, la coma se utiliza como delimitador.

Después de esto, se forma un par utilizando un registro en el índice 7 de la matriz 'SingleCountryData' y un valor '1'.

output.collect(new Text(SingleCountryData[7]), one);

Elegimos el registro en el índice 7 porque necesitamos datos de país y se encuentra en el índice 7 del array 'SingleCountryData'.

Tenga en cuenta que nuestros datos de entrada tienen el siguiente formato (donde País está en el séptimo índice, con 0 como índice inicial):

Transaction_date,Product,Price,Payment_Type,Name,City,State,Country,Account_Created,Last_Login,Latitude,Longitude

La salida del mapeador es, de nuevo, un par clave-valor que se emite utilizando el método 'collect()' de 'OutputCollector'.

Explicación de la clase SalesCountryReducer

En esta sección, comprenderemos la implementación de la clase SalesCountryReducer.

1. Comenzamos especificando el nombre del paquete para nuestra clase. SalesCountry es el nombre de nuestro paquete. Tenga en cuenta que el resultado de la compilación, SalesCountryReducer.class, se guardará en un directorio con el mismo nombre de paquete: SalesCountry.

Seguido de esto, importamos paquetes de biblioteca.

La siguiente captura de pantalla muestra una implementación de la clase SalesCountryReducer.

Vista del editor de la implementación completa de la clase SalesCountryReducer

Code Explicación:

1. Definición de clase SalesCountryReducer

public class SalesCountryReducer extends MapReduceBase implements Reducer<Text, IntWritable, Text, IntWritable> {

Aquí, los dos primeros tipos de datos, 'Texto' e 'IntWritable', son el tipo de datos de la clave-valor de entrada para el reductor.

La salida del mapeador tiene la forma de , Esta salida del mapeador se convierte en la entrada del reductor. Por lo tanto, para que coincida con su tipo de datos, se utilizan Text e IntWritable como tipos de datos.

Los dos últimos tipos de datos, 'Text' e 'IntWritable', son el tipo de datos de salida generados por el reductor en forma de par clave-valor.

Cada clase reductora debe heredar de la clase MapReduceBase y debe implementar la interfaz Reducer.

2. Definición de la función 'reducir'

public void reduce( Text t_key,
             Iterator<IntWritable> values,                           
             OutputCollector<Text,IntWritable> output,
             Reporter reporter) throws IOException {

La entrada para el método reduce() es una clave con una lista de múltiples valores.

Por ejemplo, en nuestro caso, será-

, , ,, , .

Esto se le da al reductor como

Por lo tanto, para aceptar argumentos de esta forma, se utilizan los dos primeros tipos de datos, a saber, Texto e Iterador. . El texto es un tipo de datos de clave e iterador. es un tipo de dato para una lista de valores para esa clave.

El siguiente argumento es de tipo OutputCollector. que recoge la salida de la fase reductora.

El método reduce() comienza copiando el valor de la clave e inicializando el contador de frecuencia a 0.

Text key = t_key;
int frequencyForCountry = 0;

Luego, utilizando un bucle 'while', iteramos a través de la lista de valores asociados con la clave y calculamos la frecuencia final sumando todos los valores.

 while (values.hasNext()) {
            // replace type of value with the actual type of our value
            IntWritable value = (IntWritable) values.next();
            frequencyForCountry += value.get();
            
        }

Ahora, enviamos el resultado al colector de salida en forma de clave y recuento de frecuencia obtenido.

El siguiente código hace esto:

output.collect(key, new IntWritable(frequencyForCountry));

Explicación de la clase SalesCountryDriver

En esta sección, comprenderemos la implementación de la clase SalesCountryDriver.

1. Comenzamos especificando un nombre de paquete para nuestra clase. SalesCountry es el nombre de nuestro paquete. Tenga en cuenta que el resultado de la compilación, SalesCountryDriver.class, se guardará en un directorio con el mismo nombre de paquete: SalesCountry.

Aquí hay una línea que especifica el nombre del paquete seguido de un código para importar paquetes de biblioteca.

Declaración de embalaje y declaraciones de importación de SalesCountryDriver

2. Defina una clase de controlador que creará un nuevo trabajo de cliente, un objeto de configuración y anunciará las clases Mapper y Reducer.

La clase de controlador es responsable de configurar nuestro trabajo MapReduce para que se ejecute HadoopEn esta clase, especificamos el nombre del trabajo, el tipo de datos de entrada/salida y los nombres de las clases mapper y reducer.

Código del controlador que establece el nombre del trabajo, las clases de clave y valor, y los formatos.

3. En el siguiente fragmento de código, configuramos directorios de entrada y salida que se utilizan para consumir el conjunto de datos de entrada y producir resultados, respectivamente.

arg[0] y arg[1] son ​​los argumentos de la línea de comandos que se pasan con un comando dado en la práctica de MapReduce, es decir,

$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales

Código del controlador que pasa las rutas de entrada y salida desde la línea de comandos.

4. Activar nuestro trabajo

El siguiente código inicia la ejecución del trabajo MapReduce.

try {
    // Run the job 
    JobClient.runJob(job_conf);
} catch (Exception e) {
    e.printStackTrace();
}

Preguntas Frecuentes

Las tres clases importan org.apache.hadoop.mapred, la API original de MapReduce. Hadoop 3.x todavía la incluye y la ejecuta, pero el trabajo nuevo normalmente se escribe con org.apache.hadoop.mapreduce, que reemplaza JobConf y JobClient con Configuration y Job.

Los modelos entrenados con el historial de trabajos anteriores predicen el tiempo de ejecución, sugieren tamaños de división y cantidad de reductores, y detectan desviaciones antes de que finalice la ejecución. También señalan los trabajos cuyos contadores de registros perdidos o tareas fallidas se desvían del rango habitual.

Copilot gestiona correctamente las partes repetitivas: firmas de clases, genéricos, importaciones y llamadas de configuración del controlador. Sin embargo, la lógica de negocio, como por ejemplo qué columna contiene el país, aún debe ser verificada por un desarrollador con respecto al esquema real.

Hadoop se niega a escribir en un directorio de salida existente para que los resultados finales nunca se sobrescriban. Elimine mapreduce_output_sales recursivamente con hdfs dfs -rm -r, o pase una ruta de salida diferente en la siguiente ejecución.

La herramienta jar ignora una línea final del manifiesto que no tiene un terminador de línea, por lo que Main-Class se omite silenciosamente. El archivo jar se compila sin errores, pero falla en tiempo de ejecución porque no se registra ninguna clase principal.

Sí. El ejemplo especifica la versión 2.2.0 en cada nombre de archivo JAR. En otra versión, sustituya esa versión o simplemente utilice la salida de `hadoop classpath`, que muestra todos los archivos JAR que necesita la distribución instalada.

No. Una instalación pseudodistribuida de un solo nodo lo ejecuta sin modificaciones, ya que los comandos solo asumen que HDFS y YARN están iniciados. El mismo archivo JAR se envía a un clúster multinodo real sin ningún cambio en el código.

Cambie el índice de la matriz en el mapeador de 7 a 5, ya que Ciudad es la sexta columna de SalesJan2009.csv. Recompile, reconstruya el archivo jar y ejecute la tarea en un directorio de salida nuevo.

Resumir este post con: