Hadoop MapReduce 示例:第一个 Java 程序与 Code
在本教程中,您将学习如何使用 Hadoop 和 MapReduce 示例。使用的输入数据是 2009 年 XNUMX 月销售额.csv它包含销售相关信息,例如产品名称、价格、付款方式以及客户所在的城市和国家/地区。目标是找出每个国家/地区的产品销量。
第一个 Hadoop MapReduce 程序
现在在这个 MapReduce 教程,我们将创建我们的第一个 Java MapReduce 程序:
下面的屏幕截图显示了 2009 年 1 月的原始销售数据,其中每一行代表一笔交易,国家/地区位于第八列,以逗号分隔。
请确保您已安装 Hadoop。在开始实际操作之前,请将用户更改为“hduser”(配置 Hadoop 时使用的 ID——您可以切换到您自己的 Hadoop 配置过程中使用的用户 ID)。
su - hduser_
提示信息变为 hduser 帐户,如下所示。
步骤 1)创建项目目录和源文件
创建一个名为 MapReduceTutorial 的新目录,如下面的 MapReduce 示例所示。
sudo mkdir MapReduceTutorial
授予权限
sudo chmod -R 777 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(); } } }
该存档文件会展开成相同的三个源文件,如下所示。
检查所有这些文件的文件权限
如果缺少“读取”权限,请授予该权限:
步骤 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/*"
导出的类路径会回显到提示符处。
步骤 3)编译 Java 档
编译 Java 这些文件位于 Final-MapReduceHandsOn 目录中。它们的类文件将被放入包目录中。
javac -d . SalesMapper.java SalesCountryReducer.java SalesCountryDriver.java
可以忽略此警告——它只是报告 mapred API 已弃用。
此编译过程将在当前目录中创建一个以指定的包名称命名的目录。 Java 源文件(在本例中为 SalesCountry),并将所有已编译的类文件放入其中。
步骤 4)创建清单文件
创建一个名为 Manifest.txt 的新文件
sudo gedit Manifest.txt
在其中添加以下行:
Main-Class: SalesCountry.SalesCountryDriver
SalesCountry.SalesCountryDriver 是主类的名称。请注意,您需要在本行末尾按回车键。
第五步)将这些课程打包到一个罐子里。
创建 Jar 文件
jar cfm ProductSalePerCountry.jar Manifest.txt SalesCountry/*.class
检查 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 /
我们可以安全地忽略这个警告。
验证文件是否确实被复制。
$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce
步骤 8)运行 MapReduce 作业
运行 MapReduce 作业
$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales
这将在上创建一个名为 mapreduce_output_sales 的输出目录 高密度文件系统。此目录的内容将是一个包含每个国家/地区产品销售情况的文件。
步骤 9)阅读结果
结果可以通过命令行界面查看,如下所示:
$HADOOP_HOME/bin/hdfs dfs -cat /mapreduce_output_sales/part-00000
也可以通过 Web 界面查看结果:
可选 http://localhost:50070/ 在 Web 浏览器中。在 Hadoop 3.x 中,NameNode Web UI 已移至端口。 9870所以使用 http://localhost:9870/ 相反。
现在选择“浏览文件系统”,然后导航到 /mapreduce_output_sales 目录。
打开 part-r-00000
SalesMapper 类的解释
整个作业流程如下:接下来的三个部分将详细介绍每个类实际执行的操作。
在本节中,我们将了解 SalesMapper 类的实现。
1. 首先,我们需要为类指定一个包名。SalesCountry 就是我们的包名。请注意,编译后的输出文件 SalesMapper.class 将位于以该包名命名的目录中:SalesCountry。
接下来,我们导入库包。
下图展示了 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 类的一个实现——
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。
这是指定包名称的一行,后跟用于导入库包的代码。
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(); }





















