Hadoop MapReduce 连接和计数器示例
⚡ 智能摘要
MapReduce 连接将两个大型数据集基于共享键合并,可以在 mapper 内部或 reducer 内部进行;而 MapReduce 计数器收集有关作业的统计信息,以便可以衡量错误记录而不是猜测。

MapReduce 中的 Join 是什么?
MapReduce 连接操作用于合并两个大型数据集。然而,这个过程需要编写大量代码来执行实际的连接操作。连接两个数据集首先要比较它们的大小。如果一个数据集比另一个数据集小,则较小的数据集会被分发到集群中的每个数据节点。
一旦加入 映射简化 如果采用分布式架构,则 Mapper 或 Reducer 会使用较小的数据集从较大的数据集中查找匹配的记录,然后将这些记录组合起来形成输出记录。
连接类型
根据实际执行连接的位置,Hadoop 中的连接分为两种类型。
- 地图侧连接 — 当连接操作由映射器执行时,称为映射端连接。在这种类型中,连接操作在映射函数实际使用数据之前执行。每个映射的输入必须以分区的形式存在,并且必须按排序顺序排列。此外,分区的数量必须相等,并且必须按连接键排序。
- 缩减侧连接 — 当连接操作由 reducer 执行时,称为 reduce 端连接。这种连接不需要数据集是结构化的(或分区的)。在这里,map 端处理会生成连接键以及两个表中对应的元组。经过此处理,所有具有相同连接键的元组都会进入同一个 reducer,然后由该 reducer 连接具有相同连接键的记录。
下图描述了 Hadoop 中连接的总体过程流程。

明确了两种变体之后,下一节将介绍对两个小部门文件进行归约侧连接。
如何连接两个数据集:MapReduce 示例
有两个不同的文件(如下所示),其中包含两组数据。两个文件中都包含相同的键值 Dept_ID。目标是使用 MapReduce Join 将这两个文件合并。
输入: 输入数据集为两个txt文件,分别是DeptName.txt和DeptStrength.txt。
确保你有 Hadoop的 已安装。在开始 MapReduce Join 示例的实际流程之前,请将用户更改为“hduser”(Hadoop 配置期间使用的 ID,您可以切换到 Hadoop 配置期间使用的用户 ID)。
su - hduser_
提示符变为 Hadoop 帐户,如下所示。
步骤1) 将 zip 文件复制到你选择的位置
步骤2) 解压 Zip 文件
sudo tar -xvf MapReduceJoin.tar.gz
前tractar 解压归档文件时,文件名会滚动显示。
步骤3) 进入目录 MapReduceJoin/
cd MapReduceJoin/
步骤4) 启动 Hadoop
$HADOOP_HOME/sbin/start-dfs.sh
$HADOOP_HOME/sbin/start-yarn.sh
这两个脚本都会打印出它们启动的守护进程。
步骤5) DeptStrength.txt 和 DeptName.txt 是此 MapReduce Join 示例程序使用的输入文件。
需要将这些文件复制到 高密度文件系统 使用以下命令-
$HADOOP_HOME/bin/hdfs dfs -copyFromLocal DeptStrength.txt DeptName.txt /
步骤6) 使用以下命令运行程序 -
$HADOOP_HOME/bin/hadoop jar MapReduceJoin.jar MapReduceJoin/JoinDriver/DeptStrength.txt /DeptName.txt /output_mapreducejoin
首先会回显该命令,然后作业会在控制台上报告其进度。
步骤7) 执行完成后,输出文件(名为“part-00000”)将存储在 HDFS 的 /output_mapreducejoin 目录中。
可以使用命令行界面查看结果
$HADOOP_HOME/bin/hdfs dfs -cat /output_mapreducejoin/part-00000
也可以通过 Web 界面查看结果:
现在选择“浏览文件系统”,然后导航到 /output_mapreducejoin 目录。
打开 part-r-00000
结果显示
注意: 请注意,下次运行此程序之前,您需要删除输出目录 /output_mapreducejoin
$HADOOP_HOME/bin/hdfs dfs -rm -r /output_mapreducejoin
另一种方法是对输出目录使用不同的名称。
连接操作可以告诉你数据的结构。接下来要介绍的计数器操作,则可以告诉你生成数据的作业的运行情况。
MapReduce 中的 Counter 是什么?
MapReduce 中的计数器是一种用于收集和衡量 MapReduce 作业和事件统计信息的机制。计数器用于记录…… tracMapReduce 中各种作业统计信息(例如操作次数和操作进度)的计数器用于 MapReduce 中的问题诊断。
Hadoop 计数器类似于将日志消息放入 map 或 Reduce 代码中。此信息可能有助于诊断 MapReduce 作业处理中的问题。
通常,Hadoop 中的这些计数器是在程序(map 或 reduce)中定义的,并在程序执行期间,当发生特定事件或条件(特定于该计数器)时递增。Hadoop 计数器的一个非常好的应用是…… trac从输入数据集中提取 k 条有效记录和 k 条无效记录。
MapReduce 计数器的类型
MapReduce计数器基本上有两种类型。
- Hadoop内置计数器: 每个作业都有一些内置的 Hadoop 计数器。以下是内置计数器组:
- MapReduce 任务计数器 — 在执行期间收集特定任务信息(例如,输入记录的数量)。
- 文件系统计数器 — 收集诸如任务读取或写入的字节数之类的信息。
- 文件输入格式计数器 — 收集通过 FileInputFormat 读取的若干字节的信息。
- 文件输出格式计数器 — 收集通过 FileOutputFormat 写入的若干字节的信息。
- 工作计数器 — 这些计数器记录整个作业的统计数据,例如为一项作业启动的任务数量。
- 用户自定义计数器: 除了内置计数器之外,用户还可以使用编程语言提供的类似功能定义自己的计数器。例如,在 Java“枚举”用于定义用户自定义计数器。
💡 版本说明: 工作计数器由工作部门维护。Trac在 MRv1 下,ker 属于 MapReduce ApplicationMaster。在 YARN 上,该角色属于 MapReduce ApplicationMaster,因此计数器名称得以保留,但报告它们的组件已发生变化。
一个作业不能声明无限数量的计数器。 mapreduce.job.counters.max 默认情况下,此设置将每个作业的总请求数限制为 120,如果作业声明的请求数超过 120,则作业会失败并返回错误信息。 LimitExceededException因此,计数器是用来统计少量聚合信号的,而不是用来统计每个按键的计数。
计数器示例
这是一个使用计数器来统计缺失值和无效值数量的 MapClass 示例。本教程中使用的输入数据文件是 CSV 文件 SalesJan2009.csv。
public static class MapClass extends MapReduceBase implements Mapper<LongWritable, Text, Text, Text> { static enum SalesCounters { MISSING, INVALID }; public void map ( LongWritable key, Text value, OutputCollector<Text, Text> output, Reporter reporter) throws IOException { //Input string is split using ',' and stored in 'fields' array String fields[] = value.toString().split(",", -20); //Value at 4th index is country. It is stored in 'country' variable String country = fields[4]; //Value at 8th index is sales data. It is stored in 'sales' variable String sales = fields[8]; if (country.length() == 0) { reporter.incrCounter(SalesCounters.MISSING, 1); } else if (sales.startsWith("\"")) { reporter.incrCounter(SalesCounters.INVALID, 1); } else { output.collect(new Text(country), new Text(sales + ",1")); } } }
以上代码片段展示了 Hadoop MapReduce 中计数器的一个实现示例。
在这里, 销售柜台 是使用“定义的计数器”枚举它用于统计缺失和无效的输入记录。
在代码片段中,如果“国家如果字段长度为零,则其值缺失,因此相应的计数器 SalesCounters.MISSING 递增。
接下来,如果“销售如果字段以“ 开头,则该记录被视为无效。这将通过递增计数器 SalesCounters.INVALID 来表示。
💡 API 说明: 上面的代码片段使用了原始版本 org.apache.hadoop.mapred API,其中 MapReduceBase, Mapper 接口, OutputCollector 和 Reporter 单独显示。当前代码是针对……编写的。 org.apache.hadoop.mapreduce其中单个 Context 替换收集器和报告器,并且计数器递增 context.getCounter(SalesCounters.MISSING).increment(1)两者的计数器概念是相同的。











