ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

MapReduce、Hive与Pig:大数据批处理核心原理与实战避坑指南

MapReduce、Hive与Pig:大数据批处理核心原理与实战避坑指南 直接聊点实在的。大数据批处理这个圈子绕不开“MapReduce / Hive / Pig”这三个名字。很多人刚接触时会被绕晕MapReduce 是 Java 代码Hive 是 SQLPig 又是一套脚本语言三者到底什么关系实际项目里该用哪一个为什么现在 Flink、Spark 满天飞还有人在提 Hive 优化小文件、自定义 UDAF这些问题如果不从底层批处理模型理清楚后面做离线数仓、数据清洗、综合实训多半会栽跟头。这篇内容我不会给你讲教科书式的概念而是按我自己从写 MapReduce 原始代码到用 Hive 跑 SQL、再用 Pig 做 ETL 的真实路径把三个东西的底层逻辑、实操要点和常见坑串一遍。适合刚入门大数据、准备做离线数仓作业、或者正在复习架构原理的读者。你能得到的不只是“怎么用”还有“为什么这么设计”。1. 先理清关系MapReduce、Hive、Pig 到底在解决什么问题1.1 一句话说清三者的分工MapReduce 是一套分布式计算框架它规定了你用什么模型去处理大规模数据先把任务拆成 Map映射阶段和 Reduce归约阶段中间依靠框架自动完成数据分组和排序。你写的是 Java 代码要自己控制 Mapper、Reducer、Partitioner、Comparator 这些组件。Hive 则是一个数据仓库工具。它把 SQL 翻译成 MapReduce后来也可以翻译成 Spark 或 Tez去执行。也就是说你写SELECT count(*) FROM table GROUP BY dtHive 负责帮你生成一堆 Mapper 和 Reducer。它省掉了你手写 Java 的麻烦本质上是“SQL 编译器 执行引擎调度器”。Pig 又是一套独立的数据流脚本语言叫 Pig Latin。它既不要求你写 Java也不用写 SQL而是用类似“加载 - 过滤 - 转换 - 存储”的流水线思维来处理数据。Pig 也会把脚本编译成 MapReduce 去跑但它的抽象层级介乎于原生 MapReduce 和 Hive 之间比手写 Java 简单比 SQL 更灵活特别适合做复杂 ETL 清洗。用一句话概括MapReduce 是“发动机”Hive 和 Pig 是两种不同的“方向盘”。发动机都是那个分布式批处理模型但你怎么跟它交互决定了你的姿势。1.2 为什么现在还要回头看这套组合很多初学者问现在离线计算都用 Spark SQL、Flink SQL谁还碰 MapReduce 和 Pig这个问题我当年也问过。实际去面试、做实训、维护老项目时会发现大量公司的离线数仓核心链路依然是 Hive因为数据量到一定规模后稳定性和语法成熟度比“新引擎”更重要。Pig 虽然热度下降但在某些电商、日志分析老系统里还有存量脚本。更重要的是理解了 MapReduce 模型你才能真正看懂 Shuffle 为什么贵、小文件为什么毒、数据倾斜为什么疼。Spark 的 Shuffle 原理就是从 MapReduce 演化过来的只是把中间结果优先放内存。如果你直接学 Spark SQL完全跳过底层遇到数据倾斜时你根本不知道去调哪个参数。所以把这套老组合吃透不是为了复古而是为了举一反三。2. MapReduce 底层批处理从代码到分布式执行2.1 MapReduce 的核心执行模型Map、Shuffle、Reduce 是怎么协作的我手写 MapReduce 程序时最痛苦的不是写 Mapper而是理解整个执行流程。一次完整的 MapReduce 作业数据流向大概是输入分片 - Map 阶段 - 溢写排序 - Shuffle 拷贝 - Merge 合并 - Reduce 阶段 - 输出但注意Map 的输出不是直接给 Reduce而是先写到本地磁盘。每个 Map 任务处理完一个分片后会产生一个输出文件这个文件内部已经按 key 进行了分区默认用哈希取模和排序。然后 Reduce 任务从各个 Map 任务的输出里把自己负责的分区数据拉取过来再做一次合并排序最后逐 key 调用 reduce 函数。这里必须解释“为什么中间要落盘”。老版本 MapReduce 的中间结果必须写到磁盘是因为设计时假设数据量极大内存根本装不下而且当时没有好的容错机制——如果中间结果只放内存某个节点挂了数据全丢没法重算。Spark 把中间结果尽量放内存所以快但它需要更精细的内存管理和容错机制。理解了这一点你就不会疑惑为什么像“分组排序”这种在 SQL 里一行代码的事在 MapReduce 里要写一堆类。在代码层面一个完整作业至少需要四个组件Mapper实现 map 方法、Reducer实现 reduce 方法、Driver设置作业参数并提交、可选的自定义 Writable 和 Comparator。你的 map 方法输入是key, value输出是newKey, newValuereduce 方法接收的是同一个 key 下面的所有 value 列表然后做聚合。2.2 最经典的分组排序案例自定义 WritableComparator 到底改了什么网上搜索“MapReduce 排序—分组排序”的实训题特别多核心点都在于“排序”到底由谁控制。我举一个最典型的场景有一批订单数据字段是用户ID、订单金额、下单时间现在要按“用户ID升序金额降序”输出。如果你直接把组合键设成TextIntWritable默认排序是“第一个字段升序第二个字段升序”显然不符合要求。这里就要用自定义 key重写 compareTo 方法。比如public static class OrderKey implements WritableComparableOrderKey { private Text userId; private DoubleWritable amount; public int compareTo(OrderKey other) { int cmp userId.compareTo(other.userId); if (cmp ! 0) return cmp; return -amount.compareTo(other.amount); // 金额降序 } }但这里有个陷阱如果分组也按这个 key 来那么同一用户的每一条订单都会被视为不同的 keyreduce 会调用很多次而不是一次拿到该用户所有订单。所以你要重写GroupingComparator让它只按 userId 分组public static class OrderGroupingComparator extends WritableComparator { protected OrderGroupingComparator() { super(OrderKey.class, true); } public int compare(WritableComparable a, WritableComparable b) { OrderKey oa (OrderKey) a; OrderKey ob (OrderKey) b; return oa.getUserId().compareTo(ob.getUserId()); } }然后在 Driver 里设置job.setGroupingComparatorClass(OrderGroupingComparator.class);这就是“分组排序”的核心排序规则由 key 的 compareTo 决定分组规则由 GroupingComparator 决定两者可以不一样。我在实训里见过很多人只重写 compareTo 而不改 GroupingComparator结果排序对了但 reduce 没拿到整组数据后面做 TopN 全乱。2.3 常见 MapReduce 编程踩坑点手写 MapReduce 最容易出问题的细节我按踩过的频率排个序第一输入输出路径传参。在本地跑通后丢到集群上最常见的错误是输出目录已经存在或者输入路径带通配符没被正确解析。建议在 Driver 里先做一次路径检查用FileSystem.exists判断存在就删除或报错。第二数据类型必须实现 Writable。你自建的 POJO 不能直接作为 key必须实现WritableComparable否则序列化和排序都会炸。尤其是嵌套类型别偷懒老老实实写 write/readFields。第三Reduce 端遍历 values 时不能保留 key 引用。框架会复用同一个 key 对象如果你把 key 直接存到集合里最后会发现所有元素都是同一个对象。解决办法是new Text(key)拷贝一份。第四小文件本身是毒药。MapReduce 默认把 128MB 作为一个分片如果你有十万个 1KB 的小文件就会启动十万个 Map 任务调度开销比计算本身还大。实战中先合并小文件或者用 Hive 的CombineHiveInputFormat不然跑一次作业要等半小时调度。3. Hive SQL 落地从建表到数据清洗3.1 建表与数据加载内外表、分区表和分桶表的选型Hive 最舒服的地方是能用 SQL 表达复杂逻辑。但“SQL 一时爽建表火葬场”也是真的。我见过太多人把 Hive 当成 MySQL 用建表时不分区随着数据日积月累查询全表扫描慢到怀疑人生。先说内外表。外部表External Table的数据路径由你指定drop 表只是删元数据文件还在内部表Managed Table由 Hive 管理数据文件drop 表会连文件一起删。对于清洗后的结果表我习惯用内部表生命周期明确对于原始日志、业务库同步过来的数据用外部表更安全避免误删源文件。分区表是离线数仓的命根子。按天分区是最常见的策略比如dt string分区。查询时加上WHERE dt 2025-04-01Hive 会直接定位到对应分区目录而不是全表扫描。但分区数不能无限膨胀如果每分钟一个分区元数据本身就要压垮 NameNode。一般我建议日活亿级以内的表按天或按小时分区足够。分桶表则用于采样、Join 优化和 Bucket Map Join。分桶列经过哈希取模落到固定数量的桶文件里两个表如果分桶列相同、桶数成倍数关系可以走 bucket join避免全量 Shuffle。但这个优化依赖hive.enforce.bucketingtrue和实际文件布局匹配新手先不用强求分区已经能解决 80% 的性能问题。建表语句模板我一般是这样的CREATE EXTERNAL TABLE ods.order_log ( user_id STRING, amount DECIMAL(10,2), order_time TIMESTAMP ) PARTITIONED BY (dt STRING) ROW FORMAT DELIMITED FIELDS TERMINATED BY \t STORED AS TEXTFILE LOCATION /data/ods/order_log;然后加载数据LOAD DATA INPATH /tmp/order_log_20250401 INTO TABLE ods.order_log PARTITION(dt2025-04-01);这里注意LOAD DATA是移动文件不是复制如果源路径和表路径在同一文件系统它会直接 move。如果你想让之前已经有的分区数据被识别需要MSCK REPAIR TABLE修复分区元数据。3.2 Hive 小文件问题为什么会产生怎么合并搜索引擎热词里“Hive 优化小文件”常年霸榜确实绕不开。小文件是指在 HDFS 上大量小于 Block 大小的文件。它们的危害不只是查询时 Map 启动数量爆炸更麻烦的是元数据膨胀NameNode 内存里每个文件都有一个条目百万小文件能吃掉几十 GB 内存。小文件产生的路径非常多上游 Spark/Flink 写入时并行度太高、多个 Reduce 任务输出后没有合并、频繁INSERT INTO附加数据、动态分区插入时每个分区写入的文件数等于 reducer 数。我遇到最典型的是 Flink 写 Hive 表时checkpoint 频繁触发每个 checkpoint 都产生一批小文件。合并小文件有两个方向。一个是治标——在查询时让 Hive 自动合并小文件作为输入SET hive.input.formatorg.apache.hadoop.hive.ql.io.CombineHiveInputFormat; SET hive.merge.mapfilestrue; SET hive.merge.mapredfilestrue; SET hive.merge.size.per.task256000000; SET hive.merge.smallfiles.avgsize16000000;这个方案能让一个 Map 任务处理多个小文件减少任务数但它不改变 HDFS 上的文件数量元数据压力还在。治本的方法是重建表把数据重新写一遍让文件数收敛。常用办法是开个临时表控制 Reduce 数量SET mapreduce.job.reduces1; CREATE TABLE tmp_order_log AS SELECT * FROM order_log WHERE dt2025-04-01;也可以直接用INSERT OVERWRITE重写原表但要小心不要中途失败导致数据丢失。更稳妥的方法是先写临时表确认无误后切换分区或表名。从源头控制更值得做。如果是 Flink 写 Hive可以配置sink.partition-commit.policy和sink.rolling-policy.file-size让文件在落盘前触发滚动而不是每个 checkpoint 都刷一个文件。如果你控制不了上游那就只能定期跑一个合并任务把一天里的小文件压缩成一个大文件。3.3 自定义 UDAF写一个简单聚合函数的完整流程有些业务聚合逻辑 Hive 内置函数不够用比如你自己要算“中位数”、“分位数”、“加权平均”就需要写 UDAFUser Defined Aggregate Function。热词里也有“Hive 自定义 UDAF函数”说明现在还在考这些。UDAF 的抽象类有两个老版的GenericUDAFEvaluator和新版的GenericUDAFResolver2。我建议直接用AbstractGenericUDAFResolver搭配GenericUDAFEvaluator。核心是实现四个阶段的 EvaluatorINIT 阶段初始化聚合上下文确定输入输出类型。ITERATE 阶段每来一行数据更新中间状态。TERMINATE_PARTIAL 阶段返回部分聚合结果用于 Map 端合并。MERGE 阶段合并多个部分聚合结果。TERMINATE 阶段输出最终结果。举个计算“平均绝对值偏差”的例子。你要对一列数值求每个值与均值的绝对差的平均值。首先需要自定义一个中间状态类AverageAbsoluteDeviationBuffer里面存数值个数、总和、绝对差累加值。但难点在于平均绝对偏差需要先知道均值而均值要等所有数据都聚合完。这种“两遍计算”的需求不能简单用单遍聚合实现。实际中我会拆成三步先查均值再用 SQL 计算绝对差最后求平均。UDAF 只适合那些能通过“部分合并”完成的算法比如求和、计数、分位数估算。如果你一定要写一个真正的 UDAF可以参考最简单的Sum实现把处理逻辑替换成自己的公式。写完编译成 jar然后在 Hive 里ADD JAR /path/to/udaf.jar; CREATE TEMPORARY FUNCTION my_avg_dev AS com.example.MyAvgDevUDAF; SELECT my_avg_dev(amount) FROM order_log;注意临时函数只在当前会话有效跑完就没了。如果要在生产环境用把它放到 Hive 的 auxlib 目录或者用永久函数注册。4. Pig 脚本比 Hive 更轻量的批处理语言4.1 Pig Latin 的核心算子和执行逻辑我在用 Pig 之前以为它就是个简化版 Hive用完之后才发现它和 SQL 的思路完全不同。SQL 是声明式的——你告诉系统“要什么”系统决定“怎么做”。Pig 是数据流式的——你告诉系统“数据先怎样、再怎样”更像 Unix 管道cat file | grep | sort | uniq。Pig 的核心算子就那几个LOAD、FILTER、FOREACH、GROUP、JOIN、SPLIT、STORE。每一个算子对应 MapReduce 里面的一到多个阶段。例如records LOAD hdfs:///data/logs/ USING PigStorage(,) AS (user_id:chararray, action:chararray, ts:long); filtered FILTER records BY action buy; grouped GROUP filtered BY user_id; counted FOREACH grouped GENERATE group AS user_id, COUNT(filtered) AS buy_cnt; sorted ORDER counted BY buy_cnt DESC; STORE sorted INTO /tmp/buy_count USING PigStorage(\t);这段脚本的执行计划会被 Pig 翻译成 MapReduce 作业。注意GROUP会产生一个 key 和对应的 bag包FOREACH里面可以用COUNT、SUM、TOKENIZE等内置函数。ORDER是全量排序代价很高不到最后输出结果时不要轻易用。Pig 的优势在于写复杂 ETL 时不用嵌套一堆临时表。比如你要对数据进行多重过滤、字段拆分、多路输出以脚本线的形式组织起来非常直观。而且动态执行语法可以在FOREACH里写条件表达式比如CASE这在老版本 Hive 里要绕很大一圈。4.2 一个完整的 Pig ETL 脚本示例招聘数据清洗网络热词里有一个“实验4 mapreduce综合应用案例 — 招聘数据清洗”我拿这个场景改成 Pig 脚本来展示。假设原始文件是招聘平台导出的 CSV字段包括公司名、岗位名称、薪资范围、学历要求、发布时间。清洗需求是剔除薪资为空的记录、修正学历字段的取值、按城市拆分。-- 加载原始数据 raw LOAD /data/raw/job.csv USING PigStorage(,) AS (company:chararray, job:chararray, salary:chararray, education:chararray, city:chararray, publish:chararray); -- 过滤掉薪资为空或无法解析的 with_salary FILTER raw BY salary is not null AND salary matches .*[0-9].*; -- 用 FOREACH 来解析薪资区间提取最低薪和最高薪 parsed FOREACH with_salary GENERATE company, job, FLATTEN(STRSPLIT(salary, -, 2)) AS (salary_low:chararray, salary_high:chararray), education, city, publish; -- 学历字段标准化 clean FOREACH parsed GENERATE company, job, (chararray)REPLACE(education, 本科及以上, 本科) AS edu_clean, city, publish; -- 按城市分组计算岗位数量 by_city GROUP clean BY city; city_cnt FOREACH by_city GENERATE group AS city, COUNT(clean) AS cnt; STORE clean INTO /data/clean/job_20250401 USING PigStorage(\t); STORE city_cnt INTO /data/clean/job_city_cnt USING PigStorage(\t);这个脚本里有几个皮毛操作要解释一下。FLATTEN(STRSPLIT(...))会把拆分出的字段“炸开”因为原薪资可能是8千-1万我们需要它变成两列。REPLACE是 Pig 内置的字符串替换。最后STORE可以指定多个输出路径实现一次清洗多路落地。Pig 对 schema 不是强制的即使你写错类型它也会用“字节数组”先糊弄过去直到 store 时才可能报错。我建议在 LOAD 时尽量把类型写死比如salary:chararray而不是不写类型否则后面做比较运算时会出现奇怪的 cast 异常。4.3 Pig 与 Hive 的选型对比用了两个项目后我的判断很简单如果团队里数据分析师居多他们都会 SQL那就用 Hive。如果团队里有开发倾向的大数据工程师要做复杂的逐步清洗、多路输出、自定义 transformPig 的脚本流更顺手。还有一个维度性能调优。Hive 的优化大部分靠参数比如hive.auto.convert.join、hive.exec.parallelPig 的优化在脚本结构比如尽早 FILTER 减少数据量、选择合适的 Parallel 参数。Pig 的PARALLEL 20可以控制 reducer 数量如果你在 GROUP 后面忘了设置默认 reducer 数会根据集群容量来波动比较大。我遇到过一个场景上游日志格式不规律有的行多几个字段有的行少几个字段用 SQL 写要处理 schema 漂移问题用 Pig 就很舒服因为LOAD后可以用SIZE函数判断字段数量灵活处理脏数据。这种灵活性是高阶 SQL 需要写一堆表达式才能模拟的。5. 实操中的典型问题与排查技巧5.1 Flink 写 Hive 表数据不入表的排查思路搜索引擎热词中有“flink sink hive表 数据不入表”这我太熟悉了。Flink 写 Hive 表有几个阶段先写到临时目录等触发 checkpoint 或提交策略后才通过 metastore 感知变化执行文件提交和分区添加。如果你发现“数据没入表”其实分三种情况。第一种数据已经在临时文件里只是还没发生 checkpoint。流作业默认每 60 秒一次 checkpoint如果任务才运行几秒钟你查 Hive 自然看不到。解决办法是等待或者临时调小 checkpoint 间隔。第二种分区没被添加。Hive 分区表的写入Flink 需要配置sink.partition-commit.policy。如果策略没配metastoresuccess-file文件写进去了但分区元数据没有注册。我建议配置sink.partition-commit.trigger partition-time, sink.partition-commit.policy success-file, metastore,第三种你查询太快读取到了旧快照。Hive 在读取有分区新增的表时需要重新加载分区信息你可能要重新MSCK REPAIR TABLE或者查询前设置hive.support.concurrencyfalse强制不走 snapshot。多数情况下不是数据丢只是“不可见”。5.2 删除 Hive 乱码分区的正确姿势“删除hive乱码分区”也是一个高频坑。乱码分区一般有两种来源一是上游作业把分区目录名写坏比如dt2025-04-01%0A二是动态分区插入时字段值里带了不可见字符导致 HDFS 上的分区路径变成乱码。乱码分区的麻烦在于你不能直接用ALTER TABLE ... DROP PARTITION (dt乱码)删掉因为你很难把乱码完整地打出来。我的经验是先去 Hive 元数据库通常是 MySQL里查SELECT PART_NAME FROM PARTITIONS WHERE TBL_ID (SELECT TBL_ID FROM TBLS WHERE TBL_NAMEorder_log);拿到完整的 PART_NAME 后再用精确字符串拼接删除语句。更粗暴的办法是直接删 HDFS 目录然后用MSCK REPAIR TABLE重新同步。但注意如果目录已经不在 HDFS 上重复执行MSCK会把缺失分区标记为不一致状态有时会残留元数据。所以顺序应该是先用 FS shell 确认目录实际名称再删目录再跑MSCK REPAIR最后再查一次元数据库确认剔除干净。日常预防方面我在清洗原始数据时会对分区字段强制做 trim regexp_replace比如INSERT OVERWRITE TABLE target PARTITION(clean_dt) SELECT ..., regexp_replace(trim(dt), [^0-9-], ) AS clean_dt FROM source;这能干掉 90% 的乱码来源。5.3 慢 SQL 优化速查表从 EXPLAIN 到参数调优Hive 慢 SQL 排查第一步不是调参数而是看执行计划EXPLAIN EXTENDED SELECT a.user_id, count(*) FROM order_log a JOIN user_dim b ON a.user_id b.user_id WHERE a.dt2025-04-01 GROUP BY a.user_id;看两点Join 类型MapJoin、CommonJoin、Shuffle key 的分布。如果是 CommonJoin考虑把小表加载到内存做 MapJoinSET hive.auto.convert.jointrue; SET hive.auto.convert.join.noconditionaltask.size100000000;如果 Group By 单个 key 倾斜开启二阶段聚合SET hive.groupby.skewindatatrue;如果文件数太多导致 Map 任务数量过多先用前面的合并方案。这里我再补充几个容易被忽略的点分区裁剪不生效大概率是查询中的分区字段被函数包裹比如WHERE date(dt) 2025-04-01而不是WHERE dt 2025-04-01导致 Hive 无法做裁剪只能全表扫描。列式存储TEXTFILE 该换 ORC 就换 ORC。同样的查询ORC 加 zlib 压缩扫描数据量能少一半以上IO 压力也小很多。limit 不能提前终止 Map 端读取SELECT * FROM table LIMIT 10看起来快但在 Map 阶段仍然会读取分区下所有文件只是输出时截断。真要抽样测试用TABLESAMPLE。5.4 实操经验小结数据落地前的最后一道防线最后分享一个我自己的习惯。不管用 Hive 还是 Pig最终落地到 HDFS 的表我都会做三层校验第一层用SELECT COUNT(*)对比源数据量确认没有丢失第二层抽查几个关键分区字段确认没有乱码和空值第三层检查输出目录的文件数如果小文件数量异常就先合并再进数仓。这么做看起来多花了 10 分钟能帮你在第二天被业务方问“为什么报表少了数据”时直接甩出证据链。大数据离线处理不只是“能跑”更要“可追溯”。一旦链路里某一步的输入数据发生 schema 变化越早发现代价越小。我踩过最深的坑就是觉得上层 SQL 简单忽略底层文件布局结果上游的一个小改动导致下游整整一周的数据全部错位。所以理解 MapReduce 的底层批处理、Hive 的 SQL 落地、Pig 的脚本流本质是在帮你建立一个从“数据如何存储”到“数据如何被计算”的完整视角。这个视角才是你在批处理这个领域里真正的护城河。
返回列表