简介:这是一份面向大数据初学者与Hadoop实践者的完整项目资料,围绕全国各省市酒店数据的分析与处理展开,帮助读者掌握分布式存储与MapReduce编程的核心流程。资源包共79个文件,约758KB,以Java源码与编译后的class文件为主,辅以XML配置、properties参数文件、csv数据集及part-r-00000结果文件,覆盖从代码编写到作业运行的完整链路。项目以hotel.csv为数据源,通过HDFS分布式存储,再用Java编写的Map与Reduce程序统计酒店总数、省市分布、平均房价等指标,说明文档则记录了数据清洗、任务实现与结果解读的细节。目前已有2096人学习下载,适合希望以真实案例入门Hadoop生态、理解MapReduce执行机制并积累大数据处理经验的读者参考。
1. 全国酒店数据上 Hadoop:从一堆 CSV 到能查能算的离线数仓
手里拿到一份「全国各省市酒店数据」的 CSV,几十万到几百万行不等,字段无非是酒店名、省市、地址、评分、评论数、价格、星级、开业时间这几类。用 Excel 打开直接卡死,用 pandas 单机跑一遍 group by 还能忍,但一旦要按省、市、星级、价格区间做多维交叉统计,再叠加评论数排序、评分分布,单机内存就开始告急。这时候 Hadoop 的价值就出来了:HDFS 负责把大文件切块存下来,MapReduce 或 Hive 负责把聚合逻辑分发到多台机器上跑。
这篇笔记讲的就是这条链路怎么落地:数据怎么清洗、怎么传到 HDFS、怎么用 Hive 建外部表、怎么写 SQL 做省市维度的分析、中间会遇到哪些编码和分隔符的坑。适合正在做 Hadoop 课程设计、或者第一次把业务数据往集群上搬的工程师。整套流程在伪分布式和真集群上都能跑,区别只是mapreduce.framework.name和资源参数。
2. 先想清楚数据长什么样:字段清洗与 HDFS 落地
2.1 酒店数据集的典型字段与脏数据形态
全国酒店数据这类数据集,来源通常是爬虫或者公开数据平台导出,字段结构大同小异。我一般先拿head -5和wc -l摸一遍底:
# 看前 5 行,确认表头和分隔符 head -5 hotels.csv # 统计总行数,估算数据规模 wc -l hotels.csv # 看是否有 Windows 换行符(\r\n),这是后面 Hive 查询出错的常见原因 file hotels.csv常见的脏数据有这几类:一是省市字段不统一,有的写「广东省」,有的写「广东」,有的干脆是「广东 深圳」挤在一个字段里;二是价格字段混了「¥」符号或者「起」字;三是评分字段有空值或者「暂无评分」这种文本;四是 CSV 里酒店名带英文逗号,导致列错位。这些问题不处理,后面 Hive 建表查出来的结果就是一团乱麻。
处理思路是分两步:能脚本化的用 Python 或 awk 批量清洗,清洗完再上传 HDFS。不要指望在 Hive 里用 SQL 把所有脏数据都洗干净,那样 SQL 会写得又长又难维护。
2.2 用 Python 做字段标准化与格式统一
下面这段脚本做四件事:统一省市名称、剥离价格里的非数字字符、把评分空值填成 -1、把清洗后的结果写成制表符分隔的文件(避免酒店名里的逗号干扰)。
import csv import re # 省市名称映射表,把各种写法归一到标准省名 PROVINCE_MAP = { "广东": "广东省", "广东省": "广东省", "浙江": "浙江省", "浙江省": "浙江省", "江苏": "江苏省", "江苏省": "江苏省", # 实际项目里这张表要按数据里出现的所有写法补全 } def clean_price(raw): """从 '¥388起' 这类字符串里提取数字,提取不到返回 -1""" if not raw: return -1 m = re.search(r"(\d+)", raw) return int(m.group(1)) if m else -1 def clean_score(raw): """评分字段:空值或'暂无评分'统一填 -1""" if not raw or "暂无" in raw: return -1.0 try: return float(raw) except ValueError: return -1.0 with open("hotels.csv", encoding="utf-8") as fin, \ open("hotels_clean.tsv", "w", encoding="utf-8", newline="") as fout: reader = csv.DictReader(fin) writer = csv.writer(fout, delimiter="\t") # 写出表头,字段名后面 Hive 建表要用 writer.writerow(["hotel_name", "province", "city", "price", "score", "star", "comment_cnt"]) for row in reader: province = PROVINCE_MAP.get(row["province"].strip(), row["province"].strip()) writer.writerow([ row["hotel_name"].strip(), province, row["city"].strip(), clean_price(row.get("price", "")), clean_score(row.get("score", "")), row.get("star", "").strip(), row.get("comment_cnt", "0").strip() or "0", ]) print("清洗完成,输出 hotels_clean.tsv")逻辑说明:PROVINCE_MAP是归一化的核心,实际项目里这张表要根据cut -f2 hotels.csv | sort -u的结果来补,不要凭感觉写。clean_price用正则提取第一个数字串,兼容「¥388」「388元」「388起」几种写法。输出用制表符分隔而不是逗号,是因为酒店名里带逗号的情况太常见,用逗号做分隔符迟早翻车。
参数说明:encoding="utf-8"要跟源文件编码一致,如果源文件是 GBK,这里要改成gbk,否则第一行就报UnicodeDecodeError。newline=""是 csv 模块的推荐写法,避免 Windows 下多出空行。
2.3 上传 HDFS 并确认块分布
清洗完的文件要传到 HDFS。伪分布式环境下 NameNode 默认在localhost:9000(Hadoop 3.x 常见配置),真集群换成实际地址。
# 在 HDFS 上建目录,按业务分目录是好习惯 hdfs dfs -mkdir -p /warehouse/hotel/ods/hotels # 上传清洗后的文件 hdfs dfs -put hotels_clean.tsv /warehouse/hotel/ods/hotels/ # 确认文件在 HDFS 上的块分布和副本数 hdfs fsck /warehouse/hotel/ods/hotels/hotels_clean.tsv -files -blocks # 看文件大小,估算后面 MapReduce 会起几个 map hdfs dfs -du -h /warehouse/hotel/ods/hotels/fsck的输出会告诉你这个文件被切成了几个 block。默认块大小 128MB,如果文件只有几十 MB,就是一个 block,后面 MapReduce 只会起一个 map 任务,跑起来看着像单机。这不是 bug,是数据量还没到。真要看并行效果,要么把文件搞大,要么调小dfs.blocksize。
提示:上传前先
hdfs dfs -ls确认目标目录不存在同名文件,-put遇到同名文件会直接报错退出,不会覆盖。
3. 用 Hive 建外部表:把 HDFS 上的 TSV 变成能查的表
3.1 外部表 vs 管理表:为什么酒店数据要用外部表
Hive 建表分管理表(managed table)和外部表(external table)。管理表删表的时候会把 HDFS 上的数据一起删掉,外部表只删元数据,数据还在。酒店数据这种原始数据,我一般用外部表,因为后面可能还要用 Spark 或者 MapReduce 直接读这份数据,不想因为 Hive 里 drop 一下表就把数据搞没了。
建表语句如下:
CREATE EXTERNAL TABLE IF NOT EXISTS ods_hotels ( hotel_name STRING COMMENT '酒店名称', province STRING COMMENT '省份', city STRING COMMENT '城市', price INT COMMENT '价格,-1表示缺失', score DOUBLE COMMENT '评分,-1表示缺失', star STRING COMMENT '星级', comment_cnt INT COMMENT '评论数' ) ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t' STORED AS TEXTFILE LOCATION '/warehouse/hotel/ods/hotels';逻辑说明:EXTERNAL关键字决定这是外部表。ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t'必须跟清洗时用的分隔符一致,这里清洗输出的是 TSV,所以是\t。LOCATION指向 HDFS 上已经上传好的目录,建表的时候 Hive 不会移动数据,只是把元数据登记到 metastore。
参数说明:price和comment_cnt用INT,score用DOUBLE,如果源数据里这些字段有非数字内容,查询时会返回NULL,不会报错,但结果会偏。所以清洗阶段把缺失值填成 -1 是有意义的,查询时用WHERE price > 0就能过滤掉。
建完表先验证一下:
-- 看前 10 行,确认字段没串列 SELECT * FROM ods_hotels LIMIT 10; -- 统计总行数,跟 wc -l 的结果对一下 SELECT COUNT(*) FROM ods_hotels;如果SELECT *出来的字段明显错位,八成是分隔符不对或者源文件里有残留的\r。\r的问题可以在清洗脚本里用line.replace("\r", "")处理掉。
3.2 分区表改造:按省份分区提升查询效率
全国酒店数据按省份查询是最常见的场景。如果每次查询都全表扫描,数据量大了之后很慢。Hive 分区表可以把数据按省份物理分开存,查询时只扫对应分区。
改造方式是建一张分区表,然后用INSERT ... SELECT从外部表导数据:
-- 建分区表,按省份分区 CREATE TABLE IF NOT EXISTS dw_hotels_partitioned ( hotel_name STRING, city STRING, price INT, score DOUBLE, star STRING, comment_cnt INT ) PARTITIONED BY (province STRING) STORED AS ORC; -- 开启动态分区,让 Hive 根据 province 字段自动创建分区 SET hive.exec.dynamic.partition = true; SET hive.exec.dynamic.partition.mode = nonstrict; -- 从外部表导数据到分区表 INSERT OVERWRITE TABLE dw_hotels_partitioned PARTITION (province) SELECT hotel_name, city, price, score, star, comment_cnt, province FROM ods_hotels WHERE province IS NOT NULL AND province != '';逻辑说明:PARTITIONED BY (province STRING)把省份从普通字段变成分区字段,注意SELECT里 province 要放在最后一列,跟分区字段的顺序对应。hive.exec.dynamic.partition.mode = nonstrict是必须的,否则 Hive 要求至少指定一个静态分区,动态分区插不进去。
参数说明:STORED AS ORC比 TEXTFILE 省空间、查询快,适合做中间层。如果集群开了 Tez 执行引擎,ORC 格式的收益更明显。导完之后可以SHOW PARTITIONS dw_hotels_partitioned看分区列表,确认每个省都建出来了。
注意:动态分区会按 province 的 distinct 值创建分区,如果 province 字段有大量脏值(比如空字符串、乱码),会创建出一堆没用的分区。导数据前先
SELECT DISTINCT province FROM ods_hotels看一眼。
4. 省市维度分析:几个能直接抄的 Hive SQL
4.1 按省统计酒店数量、均价、平均评分
这是最基础的多维聚合,一条 SQL 能出结果:
SELECT province, COUNT(*) AS hotel_cnt, ROUND(AVG(price), 2) AS avg_price, ROUND(AVG(score), 2) AS avg_score, SUM(comment_cnt) AS total_comments FROM dw_hotels_partitioned WHERE price > 0 AND score > 0 GROUP BY province ORDER BY hotel_cnt DESC;逻辑说明:WHERE price > 0 AND score > 0把清洗阶段填的 -1 过滤掉,避免缺失值拉低均价。ROUND保留两位小数,ORDER BY hotel_cnt DESC让酒店最多的省排前面。
参数说明:如果数据量很大,ORDER BY会触发一个 reduce 做全局排序,可能比较慢。如果只是看排名,可以改成ORDER BY hotel_cnt DESC LIMIT 20,减少数据传输。
4.2 城市价格区间分布:用 CASE WHEN 做分桶
想知道每个城市里,经济型、中端、高端的酒店各占多少,用CASE WHEN分桶:
SELECT province, city, CASE WHEN price < 200 THEN '经济型' WHEN price >= 200 AND price < 500 THEN '中端' WHEN price >= 500 AND price < 1000 THEN '高端' ELSE '豪华' END AS price_level, COUNT(*) AS cnt FROM dw_hotels_partitioned WHERE price > 0 GROUP BY province, city, CASE WHEN price < 200 THEN '经济型' WHEN price >= 200 AND price < 500 THEN '中端' WHEN price >= 500 AND price < 1000 THEN '高端' ELSE '豪华' END ORDER BY province, city, cnt DESC;逻辑说明:CASE WHEN的分桶逻辑要跟业务对齐,这里的价格阈值只是示例,实际项目里要根据数据分布来定。GROUP BY里必须把CASE WHEN表达式完整重复一遍,Hive 不支持按别名分组。
参数说明:如果分桶规则经常变,建议把阈值抽成配置,或者干脆在 ETL 阶段就多算一列price_level存到表里,查询时直接 group by 这一列,SQL 更干净。
4.3 用窗口函数排城市内酒店评分 Top N
每个城市评分最高的前 5 家酒店,用ROW_NUMBER()窗口函数:
SELECT province, city, hotel_name, score, price FROM ( SELECT province, city, hotel_name, score, price, ROW_NUMBER() OVER (PARTITION BY province, city ORDER BY score DESC, comment_cnt DESC) AS rn FROM dw_hotels_partitioned WHERE score > 0 ) t WHERE rn <= 5 ORDER BY province, city, rn;逻辑说明:PARTITION BY province, city把数据按省市分组,ORDER BY score DESC, comment_cnt DESC在组内按评分降序、评论数降序排,ROW_NUMBER()给每行编个号。外层查询过滤rn <= 5就拿到每个城市的前 5 名。
参数说明:ROW_NUMBER()遇到相同评分不会并列,如果要并列排名用RANK()或DENSE_RANK()。窗口函数在 Hive 里是支持的,但如果集群版本很老(0.11 之前),需要用UDF或者改写 SQL,这种情况现在很少见了。
5. 避坑与排查:酒店数据上 Hadoop 常见的 5 个翻车点
5.1 中文乱码:查询结果全是问号
现象:Hive 查询出来的酒店名和省市全是???或者乱码。
原因:源文件编码是 GBK,清洗脚本按 UTF-8 读,或者 Hive 建表时没指定字符集,metastore 用了默认的 latin1。
解决:清洗阶段用chardet检测源文件编码,或者直接file -i hotels.csv看。Hive 这边在建表时加TBLPROPERTIES ('charset'='utf-8'),同时确认hive-site.xml里javax.jdo.option.ConnectionURL的字符集参数带了useUnicode=true&characterEncoding=UTF-8。
5.2 字段错位:酒店名里带逗号导致列偏移
现象:SELECT *出来发现 city 字段里是酒店名的一部分,后面所有列都往右偏了。
原因:源 CSV 用逗号分隔,但酒店名里本身带逗号,csv.DictReader能处理带引号的字段,但如果源文件里逗号没被引号包起来,就会错列。
解决:清洗阶段不要用split(",")这种土办法,用csv模块的DictReader。如果源文件本身格式就不规范,先用awk -F',' 'NF!=7' hotels.csv把列数不对的行捞出来单独看。输出统一用\t分隔,从根上避开逗号问题。
5.3 动态分区报错:requires at least one static partition
现象:INSERT OVERWRITE ... PARTITION (province)执行时报FAILED: SemanticException [Error 10096]: Dynamic partition cannot be the parent of a static partition。
原因:hive.exec.dynamic.partition.mode默认是strict,要求至少有一个静态分区。
解决:执行前SET hive.exec.dynamic.partition.mode = nonstrict;。这个设置只在当前会话有效,换个会话要重新设。如果用的是 Beeline,可以在连接串里加--hiveconf hive.exec.dynamic.partition.mode=nonstrict。
5.4 小文件过多:每个分区一堆几百 KB 的文件
现象:hdfs dfs -ls /warehouse/hotel/dw_hotels_partitioned/province=广东省/看到几十个几百 KB 的小文件。
原因:动态分区插入时,每个 map 任务会为每个分区生成一个文件,map 数量多、分区多,小文件就爆炸了。
解决:导完数据后跑一次合并:
-- 合并小文件 ALTER TABLE dw_hotels_partitioned PARTITION (province='广东省') CONCATENATE;或者在建表时设置hive.merge.mapfiles = true和hive.merge.mapredfiles = true,让 Hive 在任务结束时自动合并。更彻底的办法是控制 map 数量,导数据前SET mapred.reduce.tasks = 10;限制 reduce 个数。
5.5 MapReduce 任务卡在 reduce 阶段不动
现象:任务跑到 99% 卡住,reduce 阶段一直不结束。
原因:数据倾斜。某个省的酒店数据量特别大,分到同一个 reduce 上,其他 reduce 早就跑完了,就它还在跑。
解决:先确认是不是倾斜,看 JobTracker 或者 YARN 的 UI,哪个 reduce 的输入记录数明显比别人大。如果是倾斜,可以在 SQL 里加随机前缀打散:
-- 给 province 加随机前缀,打散到多个 reduce SELECT CONCAT(province, '_', CAST(RAND() * 10 AS INT)) AS province_key, COUNT(*) FROM dw_hotels_partitioned GROUP BY CONCAT(province, '_', CAST(RAND() * 10 AS INT));这个办法会让结果多出 10 倍的行,外层再聚合一次就行。或者直接SET hive.groupby.skewindata = true;,让 Hive 自动做两阶段聚合。
6. 进阶:用 MapReduce 直接处理酒店数据与结果验证
Hive SQL 写起来快,但有些定制化的清洗逻辑或者复杂的评分计算,SQL 表达起来别扭,这时候直接写 MapReduce 更灵活。下面这个例子统计每个省的平均价格,用 Java 写一个最简单的 MapReduce 作业。
// Mapper:解析 TSV,输出 <province, price> public class HotelPriceMapper extends Mapper<LongWritable, Text, Text, IntWritable> { private Text province = new Text(); private IntWritable price = new IntWritable(); @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line = value.toString(); // 跳过表头 if (line.startsWith("hotel_name")) return; String[] fields = line.split("\t"); // 字段数不对或者价格无效的直接跳过 if (fields.length < 7) return; try { int p = Integer.parseInt(fields[3]); if (p <= 0) return; province.set(fields[1]); price.set(p); context.write(province, price); } catch (NumberFormatException e) { // 价格字段解析失败,跳过这一行 } } }// Reducer:对每个省的价格求平均 public class HotelPriceReducer extends Reducer<Text, IntWritable, Text, DoubleWritable> { private DoubleWritable result = new DoubleWritable(); @Override protected void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { long sum = 0; int count = 0; for (IntWritable val : values) { sum += val.get(); count++; } if (count > 0) { result.set((double) sum / count); context.write(key, result); } } }逻辑说明:Mapper 按\t切分,取第 4 个字段(下标 3)作为价格,第 2 个字段(下标 1)作为省份。line.startsWith("hotel_name")跳过表头,fields.length < 7过滤掉格式不对的行。Reducer 累加价格和计数,最后输出平均值。
参数说明:split("\t")如果遇到字段里本身带\t的情况会切错,但清洗阶段已经保证了输出格式,所以这里可以放心用。如果数据里可能有\t,改用split("\t", -1)保留尾部空字段。
打包提交到集群:
# 打包 mvn clean package # 提交作业 hadoop jar hotel-analysis.jar com.example.HotelPriceDriver \ /warehouse/hotel/ods/hotels/hotels_clean.tsv \ /warehouse/hotel/output/avg_price # 看结果 hdfs dfs -cat /warehouse/hotel/output/avg_price/part-r-00000 | head -20跑完之后,把 MapReduce 的结果跟 Hive SQL 的结果对一下:
-- Hive 侧的结果 SELECT province, ROUND(AVG(price), 2) AS avg_price FROM dw_hotels_partitioned WHERE price > 0 GROUP BY province ORDER BY province;两边结果应该一致。如果不一致,优先检查 MapReduce 里有没有漏掉price <= 0的过滤,或者 Hive 侧有没有把空省份算进去。这种交叉验证是我每次做完离线任务都会做的,比单纯看一个结果靠谱得多。
最后一个习惯:所有清洗脚本、Hive SQL、MapReduce 代码都放到 Git 里,按etl/、sql/、mr/分目录。酒店数据这种项目,字段和口径经常变,没有版本管理,过两周自己都忘了当时怎么算的。希望帮到你。
本文还有配套的精品资源,点击获取