Hadoop MapReduce 示例:第一个 Java 程序与 Code

⚡ 智能摘要

Hadoop MapReduce 程序编写为三种形式。 Java 类、映射器、归约器和驱动程序,经过编译、打包到 jar 中,并提交到集群以统计每个国家的销售额。

  • 🔘 资料集: SalesJan2009.csv 文件每行包含一笔交易,国家/地区位于第八列,以逗号分隔。
  • ☑️ 映射器: SalesMapper 将每一行拆分,并输出与常量值 1 配对的国家/地区。
  • 减速器: SalesCountryReducer 将每个国家/地区的到达次数相加,得到一个单一的频率计数。
  • 🧪 司机: SalesCountryDriver 命名作业,声明键和值类型,并将两个类连接在一起。
  • 🛠️ 生成: 使用 javac 编译,添加 Main-Class 清单条目,然后使用 jar cfm 打包所有内容。
  • ⚠️ 跑: 将 CSV 文件复制到 HDFS,提交 jar 文件,然后从输出目录中读取 part-00000。

Hadoop MapReduce 示例创建第一个 Java 程序

在本教程中,您将学习如何使用 Hadoop 和 MapReduce 示例。使用的输入数据是 2009 年 XNUMX 月销售额.csv它包含销售相关信息,例如产品名称、价格、付款方式以及客户所在的城市和国家/地区。目标是找出每个国家/地区的产品销量。

第一个 Hadoop MapReduce 程序

现在在这个 MapReduce 教程,我们将创建我们的第一个 Java MapReduce 程序:

下面的屏幕截图显示了 2009 年 1 月的原始销售数据,其中每一行代表一笔交易,国家/地区位于第八列,以逗号分隔。

2009 年 1 月销售数据以电子表格形式打开,显示了交易列。

请确保您已安装 Hadoop。在开始实际操作之前,请将用户更改为“hduser”(配置 Hadoop 时使用的 ID——您可以切换到您自己的 Hadoop 配置过程中使用的用户 ID)。

su - hduser_

提示信息变为 hduser 帐户,如下所示。

切换到 hduser 帐户后的终端提示符

步骤 1)创建项目目录和源文件

创建一个名为 MapReduceTutorial 的新目录,如下面的 MapReduce 示例所示。

sudo mkdir MapReduceTutorial

使用 mkdir 命令创建 MapReduceTutorial 目录

授予权限

sudo chmod -R 777 MapReduceTutorial

chmod 命令授予 MapReduceTutorial 完全权限

创建三个 Java 以下源文件位于 MapReduceTutorial 中。请注意,这三个文件都使用了较旧的版本。 org.apache.hadoop.mapred API 仍然随 Hadoop 3.x 一起提供。

销售映射器.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);
	}
}

销售国家/地区Reducer.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));
	}
}

销售国家驱动程序.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();
        }
    }
}

在此处下载文件

该存档文件会展开成相同的三个源文件,如下所示。

Extracted 存档列表 SalesMapper、SalesCountryReducer 和 SalesCountryDriver

检查所有这些文件的文件权限

长列表显示了三个文件的文件权限。 Java 源文件

如果缺少“读取”权限,请授予该权限:

chmod 命令添加读取权限 Java 源文件

步骤 2)导出 Hadoop 类路径

按照以下 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/*"

导出的类路径会回显到提示符处。

导出 Hadoop CLASSPATH 变量后的 Shell 提示符

步骤 3)编译 Java 档

编译 Java 这些文件位于 Final-MapReduceHandsOn 目录中。它们的类文件将被放入包目录中。

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

javac 输出显示编译过程中出现弃用警告

可以忽略此警告——它只是报告 mapred API 已弃用。

此编译过程将在当前目录中创建一个以指定的包名称命名的目录。 Java 源文件(在本例中为 SalesCountry),并将所有已编译的类文件放入其中。

SalesCountry 包目录包含三个已编译的类文件

步骤 4)创建清单文件

创建一个名为 Manifest.txt 的新文件

sudo gedit Manifest.txt

在其中添加以下行:

Main-Class: SalesCountry.SalesCountryDriver

gedit 窗口显示 Manifest.txt 中的 Main-Class 条目

SalesCountry.SalesCountryDriver 是主类的名称。请注意,您需要在本行末尾按回车键。

第五步)将这些课程打包到一个罐子里。

创建 Jar 文件

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

使用 jar 命令从清单构建 ProductSalePerCountry.jar 文件

检查 jar 文件是否已创建

目录列表确认已创建 ProductSalePerCountry.jar 文件

步骤 6)启动 Hadoop

启动 Hadoop

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

步骤 7)将输入文件复制到 HDFS。

将文件 SalesJan2009.csv 复制到 ~/inputMapReduce 目录。

现在使用以下命令将 ~/inputMapReduce 复制到 HDFS。

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

copyFromLocal 命令输出包含本地库警告

我们可以安全地忽略这个警告。

验证文件是否确实被复制。

$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce

hdfs dfs -ls 列出 inputMapReduce 目录中的 SalesJan2009.csv 文件

步骤 8)运行 MapReduce 作业

运行 MapReduce 作业

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

SalePerCountry 作业在集群上运行时控制台输出

这将在上创建一个名为 mapreduce_output_sales 的输出目录 高密度文件系统。此目录的内容将是一个包含每个国家/地区产品销售情况的文件。

步骤 9)阅读结果

结果可以通过命令行界面查看,如下所示:

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

使用 hdfs dfs -cat 命令打印国家/地区和销售额对

也可以通过 Web 界面查看结果:

可选 http://localhost:50070/ 在 Web 浏览器中。在 Hadoop 3.x 中,NameNode Web UI 已移至端口。 9870所以使用 http://localhost:9870/ 相反。

在浏览器中查看 NameNode Web 界面主页

现在选择“浏览文件系统”,然后导航到 /mapreduce_output_sales 目录。

HDFS 文件浏览器显示 mapreduce_output_sales 目录

打开 part-r-00000

part-r-00000 结果文件已在 HDFS 浏览器中打开

SalesMapper 类的解释

整个作业流程如下:接下来的三个部分将详细介绍每个类实际执行的操作。

在本节中,我们将了解 SalesMapper 类的实现。

1. 首先,我们需要为类指定一个包名。SalesCountry 就是我们的包名。请注意,编译后的输出文件 SalesMapper.class 将位于以该包名命名的目录中:SalesCountry。

接下来,我们导入库包。

下图展示了 SalesMapper 类的一个实现——

SalesMapper 类的完整实现的编辑器视图

样本 Code 说明:

1.SalesMapper 类定义-

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

每个映射器类都必须继承自 MapReduceBase 类,并且必须实现 Mapper 接口。

2. 定义“map”函数-

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

Mapper 类的主要部分是 'map()' 方法,它接受四个参数。

每次调用“map()”方法时,都会传递一个键值对(在本代码中为“key”和“value”)。

`map()` 方法首先将作为参数接收的输入文本拆分成多个字段。

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

这里,“,”用作分隔符。

之后,使用数组“SingleCountryData”的第 7 个索引处的记录和值“1”形成一个配对。

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

我们选择索引为 7 的记录,因为我们需要国家数据,而它位于数组“SingleCountryData”的索引 7 处。

请注意,我们的输入数据格式如下(其中国家/地区位于第 7 个索引处,起始索引为 0)——

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

mapper 的输出仍然是一个键值对,它是使用 'OutputCollector' 的 'collect()' 方法发出的。

SalesCountryReducer 类说明

在本节中,我们将了解 SalesCountryReducer 类的实现。

1. 首先,我们需要为我们的类指定一个包名。SalesCountry 就是我们的包名。请注意,编译后的输出文件 SalesCountryReducer.class 将放在一个名为 SalesCountry 的目录中。

接下来,我们导入库包。

下图展示了 SalesCountryReducer 类的一个实现——

SalesCountryReducer 类的完整实现的编辑器视图

Code 说明:

1. SalesCountryReducer 类定义-

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

这里,前两个数据类型“Text”和“IntWritable”是 reducer 的输入键值的数据类型。

映射器的输出格式为: , mapper 的输出会成为 reducer 的输入。因此,为了与数据类型保持一致,这里使用了 Text 和 IntWritable 作为数据类型。

最后两种数据类型,“Text”和“IntWritable”是 reducer 生成的输出的数据类型,以键值对的形式表示。

每个 reducer 类都必须继承自 MapReduceBase 类,并且必须实现 Reducer 接口。

2. 定义“reduce”函数

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

reduce() 方法的输入是一个键,键的值为多个值。

例如,在我们的例子中,它将是-

, , , , , 。

这是传递给归约器的。

因此,为了接受这种形式的参数,首先使用了两种数据类型,即 Text 和 Iterator。文本是一种键和迭代器的数据类型。是该键对应的值列表的数据类型。

下一个参数的类型为 OutputCollector它收集降阶阶段的输出。

reduce() 方法首先复制键值并将频率计数初始化为 0。

Text key = t_key;
int frequencyForCountry = 0;

然后,使用“while”循环,遍历与键关联的值列表,并通过将所有值相加来计算最终频率。

 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 类说明

在本节中,我们将了解 SalesCountryDriver 类的实现。

1. 首先,我们需要为类指定一个包名。SalesCountry 就是我们的包名。请注意,编译后的输出文件 SalesCountryDriver.class 将位于以该包名命名的目录中:SalesCountry。

这是指定包名称的一行,后跟用于导入库包的代码。

SalesCountryDriver 的包裹申报和进口声明

2. 定义一个驱动程序类,它将创建一个新的客户端作业、配置对象并公布 Mapper 和 Reducer 类。

驱动程序类负责设置我们的 MapReduce 作业在 Hadoop的. 在此类中,我们指定作业名称、输入/输出数据类型以及映射器和归约器类的名称。

驱动程序代码设置作业名称、键和值类以及格式

3. 在下面的代码片段中,我们设置了输入和输出目录,分别用于使用输入数据集和产生输出。

arg[0] 和 arg[1] 是在 MapReduce 实践操作中,通过命令传递的命令行参数,即:

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

驱动程序代码通过命令行传递输入和输出路径

4. 触发我们的工作

以下代码开始执行 MapReduce 作业——

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

常见问题

这三个类都导入了 org.apache.hadoop.mapred,即最初的 MapReduce API。Hadoop 3.x 仍然提供并运行它,但新的工作通常是基于 org.apache.hadoop.mapreduce 编写的,后者用 Configuration 和 Job 取代了 JobConf 和 JobClient。

基于过往作业历史记录训练的模型可以预测运行时间,建议拆分大小和归约器数量,并在运行结束前发现偏差。它们还会标记出溢出记录或失败任务计数器超出正常范围的作业。

Copilot很好地完成了重复性工作:类签名、泛型、导入和驱动程序配置调用。但业务逻辑,例如哪个列存储国家/地区信息,仍然需要开发人员对照实际模式进行检查。

Hadoop 不允许写入已存在的输出目录,因此已完成的结果永远不会被覆盖。可以使用 `hdfs dfs -rm -r` 递归删除 `mapreduce_output_sales` 目录,或者在下次运行时指定不同的输出路径。

jar 工具会忽略清单文件中缺少换行符的最后一行,因此 Main-Class 会被静默丢弃。之后 jar 文件构建过程没有报错,但运行时会失败,因为没有记录主类。

是的。示例中每个 jar 包名称都锁定了 2.2.0 版本。在其他版本中,请替换为相应的版本,或者直接使用 hadoop classpath 的输出,该命令会打印出已安装发行版所需的所有 jar 包。

不。单节点伪分布式安装无需修改即可运行,因为命令仅假定 HDFS 和 YARN 已启动。同样的 jar 包无需任何代码更改即可提交到真正的多节点集群。

将映射器中的数组索引从 7 改为 5,因为 City 是 SalesJan2009.csv 的第六列。重新编译,重建 jar 文件,并对新的输出目录运行作业。

总结一下这篇文章: