Hadoop MapReduce 连接和计数器示例

⚡ 智能摘要

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

  • 🔘 加入基本步骤: 两个数据集中较小的那个会被分发到每个数据节点,并用作查找端。
  • ☑️ 地图侧连接: 在 map 函数运行之前,需要对每个输入进行分区、均等分割,并按连接键排序。
  • 缩径侧连接: 无需分区,因为每个共享连接键的元组都会落入同一个 reducer 中。
  • 🧪 示例: DeptName.txt 和 DeptStrength.txt 被复制到 HDFS 中,并通过打包的 jar 文件按 Dept_ID 连接起来。
  • 🛠️ 计数器类型: 每个作业都附带五个内置计数器组,用户自定义计数器则声明为: Java 枚举。
  • ⚠️ 反制措施: 对每个缺失或无效的记录递增计数器,可以将数据质量问题转化为作业报告上的一个数字。

Hadoop MapReduce join 和 counter 教程及示例

MapReduce 中的 Join 是什么?

MapReduce 连接操作用于合并两个大型数据集。然而,这个过程需要编写大量代码来执行实际的连接操作。连接两个数据集首先要比较它们的大小。如果一个数据集比另一个数据集小,则较小的数据集会被分发到集群中的每个数据节点。

一旦加入 映射简化 如果采用分布式架构,则 Mapper 或 Reducer 会使用较小的数据集从较大的数据集中查找匹配的记录,然后将这些记录组合起来形成输出记录。

连接类型

根据实际执行连接的位置,Hadoop 中的连接分为两种类型。

  1. 地图侧连接 — 当连接操作由映射器执行时,称为映射端连接。在这种类型中,连接操作在映射函数实际使用数据之前执行。每个映射的输入必须以分区的形式存在,并且必须按排序顺序排列。此外,分区的数量必须相等,并且必须按连接键排序。
  2. 缩减侧连接 — 当连接操作由 reducer 执行时,称为 reduce 端连接。这种连接不需要数据集是结构化的(或分区的)。在这里,map 端处理会生成连接键以及两个表中对应的元组。经过此处理,所有具有相同连接键的元组都会进入同一个 reducer,然后由该 reducer 连接具有相同连接键的记录。

下图描述了 Hadoop 中连接的总体过程流程。

Hadoop中map端连接与reduce端连接的流程图比较
Hadoop MapReduce 中的连接类型

明确了两种变体之后,下一节将介绍对两个小部门文件进行归约侧连接。

如何连接两个数据集:MapReduce 示例

有两个不同的文件(如下所示),其中包含两组数据。两个文件中都包含相同的键值 Dept_ID。目标是使用 MapReduce Join 将这两个文件合并。

首先输入的文件列出了部门 ID 和部门名称。

文件1
第二个输入文件列出了部门 ID 以及部门实力值。

文件2

输入: 输入数据集为两个txt文件,分别是DeptName.txt和DeptStrength.txt。

从这里下载输入文件

确保你有 Hadoop的 已安装。在开始 MapReduce Join 示例的实际流程之前,请将用户更改为“hduser”(Hadoop 配置期间使用的 ID,您可以切换到 Hadoop 配置期间使用的用户 ID)。

su - hduser_

提示符变为 Hadoop 帐户,如下所示。

使用 su 命令切换到 hduser 帐户后,终端会显示以下内容。

步骤1) 将 zip 文件复制到你选择的位置

已将下载的 MapReduceJoin 归档文件放置在选定的工作目录中

步骤2) 解压 Zip 文件

sudo tar -xvf MapReduceJoin.tar.gz

前tractar 解压归档文件时,文件名会滚动显示。

控制台文件列表示例trac来自 MapReduceJoin.tar.gz

步骤3) 进入目录 MapReduceJoin/

cd MapReduceJoin/

切换到 MapReduceJoin 目录后的 Shell 提示符

步骤4) 启动 Hadoop

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

这两个脚本都会打印出它们启动的守护进程。

来自 HDFS 和 YARN 守护进程脚本的启动消息

步骤5) DeptStrength.txt 和 DeptName.txt 是此 MapReduce Join 示例程序使用的输入文件。

需要将这些文件复制到 高密度文件系统 使用以下命令-

$HADOOP_HOME/bin/hdfs dfs -copyFromLocal DeptStrength.txt DeptName.txt /

两个输入文本文件都被复制到了 HDFS 根目录中。

步骤6) 使用以下命令运行程序 -

$HADOOP_HOME/bin/hadoop jar MapReduceJoin.jar MapReduceJoin/JoinDriver/DeptStrength.txt /DeptName.txt /output_mapreducejoin

首先会回显该命令,然后作业会在控制台上报告其进度。

使用命令行启动打包好的 MapReduceJoin jar 文件

控制台输出 trac掌握 MapReduce 连接作业的进度

步骤7) 执行完成后,输出文件(名为“part-00000”)将存储在 HDFS 的 /output_mapreducejoin 目录中。

可以使用命令行界面查看结果

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

使用 cat 命令从 HDFS 打印合并部门记录

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

Hadoop Web 界面登录页面用于访问文件系统浏览器

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

浏览 HDFS 文件系统视图,找到 output_mapreducejoin 目录。

打开 part-r-00000

在浏览器视图中选择 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计数器基本上有两种类型。

  1. Hadoop内置计数器: 每个作业都有一些内置的 Hadoop 计数器。以下是内置计数器组:
    • MapReduce 任务计数器 — 在执行期间收集特定任务信息(例如,输入记录的数量)。
    • 文件系统计数器 — 收集诸如任务读取或写入的字节数之类的信息。
    • 文件输入格式计数器 — 收集通过 FileInputFormat 读取的若干字节的信息。
    • 文件输出格式计数器 — 收集通过 FileOutputFormat 写入的若干字节的信息。
    • 工作计数器 — 这些计数器记录整个作业的统计数据,例如为一项作业启动的任务数量。
  2. 用户自定义计数器: 除了内置计数器之外,用户还可以使用编程语言提供的类似功能定义自己的计数器。例如,在 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,其中 MapReduceBaseMapper 接口, OutputCollectorReporter 单独显示。当前代码是针对……编写的。 org.apache.hadoop.mapreduce其中单个 Context 替换收集器和报告器,并且计数器递增 context.getCounter(SalesCounters.MISSING).increment(1)两者的计数器概念是相同的。

常见问题

当数据一侧足够小,可以驻留在所有节点的内存中时,选择 mapper 版本,因为它完全跳过了 shuffle 操作。当数据两侧都很大或数据未排序时,选择 reducer 版本,并接受额外的网络开销。

模型通过学习历史作业数据来预测运行时间,推荐拆分大小和归约器数量,并检测计数器值的偏差。它们还会标记出溢出记录或失败任务计数器超出该流水线正常范围的作业。

Copilot 生成的 mapper 和 reducer 框架看起来还不错,但它在一个类中随意混用了旧的 mapred 包和较新的 mapreduce 包,这会导致编译失败。在信任其逻辑之前,请务必修复导入语句和方法签名。

这种机制会在任务开始前将较小的文件发送到每个节点。然后,每个映射器会将该副本加载到哈希映射中,并在本地查找匹配项,这使得映射器端的连接成为可能。

当作业完成时,它们会打印在控制台摘要中,在作业历史记录和资源管理器网页中公开,并且可以从作业对象中以编程方式读取,因此驱动程序可以断言它们并判定运行失败。

一个单独的 Context 对象。它承担了 OutputCollector 和 Reporter 过去需要分配的工作,因此输出是通过传递给 map 方法的同一个句柄写入的,计数器也是通过同一个句柄递增的。

对于大多数报道工作而言,答案是否定的。 HiveQL 加入 最终只需几行代码就能编译出相同的混洗合并模式。当合并逻辑无法用 SQL 子句表达时,手动编写代码是值得的。

Hadoop 拒绝向已存在的输出目录写入数据,这可以防止已完成的结果被覆盖。请先递归删除该目录,或者在下次运行时指定不同的输出路径。

总结一下这篇文章: