ARTICLE DETAIL

资讯详情

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

基于Hadoop的全国酒店数据离线数仓实战:从CSV到Hive聚合分析

基于Hadoop的全国酒店数据离线数仓实战:从CSV到Hive聚合分析 简介这是一份面向大数据初学者与Hadoop实践者的完整项目资料围绕全国各省市酒店数据的分析与处理展开帮助读者掌握分布式存储与MapReduce编程的核心流程。资源包共79个文件约758KB以Java源码与编译后的class文件为主辅以XML配置、properties参数文件、csv数据源及part-r-00000结果输出文件结构上覆盖源码、测试、配置与运行产物便于直接编译提交作业。项目以hotel.csv为数据源通过HDFS分布式存储再用Java编写的Map与Reduce程序完成酒店总数、省市分布、平均房价等统计说明文档则记录了数据清洗、任务实现与结果解读的完整思路。目前已有2096人学习下载适合作为课程设计、大数据实验或自学练手的参考案例读者可借此理解Hadoop集群参数配置、作业提交脚本编写及MapReduce调试排错方法快速积累海量数据处理经验。1. 全国酒店数据上 Hadoop从一堆 CSV 到能查能算的离线数仓手头拿到一份全国各省市酒店数据几十万到几百万行不等字段无非是酒店名称、省份、城市、地址、星级、评分、评论数、价格区间这些。用 Excel 打开卡到怀疑人生用 pandas 跑一遍 groupby 内存直接爆掉这时候就该把 Hadoop 搬出来了。这个标题讲的就是把这份酒店明细数据落到 HDFS用 MapReduce 或 Hive 做清洗、聚合、分析最终产出「各省酒店数量」「各城市平均评分」「星级分布」「价格带分布」这类可查可看的指标。适合谁正在做课程设计、数据分析项目或者手上真有业务数据想搭一套离线处理链路的同学。它不追求实时追求的是数据量上来之后还能稳定跑完、结果可复现。下面按「先跑通最小链路再补清洗和指标最后讲坑」的顺序展开环境以 Hadoop 伪分布式为主集群搭建同理。2. 酒店数据建模与 Hadoop 环境选型先想清楚存什么、怎么存2.1 酒店明细字段拆解与分层设计拿到原始 CSV 别急着往 HDFS 扔先做字段盘点和分层规划。典型酒店数据字段大致是这几类维度字段省份、城市、区县、酒店名称、星级、度量字段评分、评论数、最低价、最高价、文本字段地址、标签、简介。分析目标决定了你要保留哪些、派生哪些。我一般按三层来放分层目录内容格式ODS 原始层/warehouse/ods/hotel_raw原样导入的 CSV文本DWD 明细层/warehouse/dwd/hotel_detail清洗、去重、类型规整后ParquetDWS 汇总层/warehouse/dws/hotel_agg按省/市/星级聚合结果Parquet为什么要分层因为原始数据一定有脏的——价格写成「300起」、评分是空字符串、省份有「广东省」和「广东」两种写法。分层之后清洗逻辑只作用在 ODS→DWD 这一段后面所有分析都基于干净的 DWD出问题好回溯。Parquet 列式存储对后面 Hive 聚合查询的扫描量优化很明显酒店数据字段多、分析时往往只取几列列存优势直接体现。分区键的选择也有讲究。酒店数据天然按省份、城市分布但省份只有三十多个做一级分区粒度太粗按城市分区又太碎。常见做法是按dt数据日期做分区如果是一次性快照数据就按province做分区单分区数据量控制在几百 MB 级别比较舒服。2.2 伪分布式还是集群按数据量选数据量在百万行以内、字段十几个单机伪分布式完全够用别一上来就折腾多节点。伪分布式就是 NameNode、DataNode、ResourceManager、NodeManager 全在一台机器上能完整跑通 HDFS 读写和 YARN 调度用来验证逻辑、写课程设计绰绰有余。判断标准很简单原始 CSV 小于 1GB伪分布式1GB 到几十 GB三节点起步再往上才考虑更大集群。酒店数据这种体量绝大多数场景伪分布式就能交付。环境准备的核心几步以 Linux 为例Windows 下用 IDEA 连远程集群同理# 1. 确认 JDKHadoop 3.x 建议 JDK 8 java -version # 2. 配置免密登录伪分布式也需要 ssh-keygen -t rsa -P -f ~/.ssh/id_rsa cat ~/.ssh/id_rsa.pub ~/.ssh/authorized_keys chmod 600 ~/.ssh/authorized_keys # 3. 解压 Hadoop 并配置环境变量 tar -zxvf hadoop-3.x.tar.gz -C /opt/ echo export HADOOP_HOME/opt/hadoop ~/.bashrc echo export PATH$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin ~/.bashrc source ~/.bashrc参数说明ssh-keygen的-P 表示空密码方便脚本免交互chmod 600是 SSH 对 authorized_keys 权限的硬要求权限不对直接拒绝登录这是新手第一个翻车点。核心配置文件要改四个core-site.xml指定fs.defaultFS为hdfs://localhost:9000hdfs-site.xml把dfs.replication设为 1伪分布式只有一块盘设 3 会一直报副本不足mapred-site.xml把mapreduce.framework.name设为yarnyarn-site.xml配yarn.nodemanager.aux-services为mapreduce_shuffle。改完格式化一次 NameNodehdfs namenode -format start-dfs.sh start-yarn.sh jps # 应看到 NameNode、DataNode、ResourceManager、NodeManagerjps少进程是排查起点少了 NameNode 看日志logs/hadoop-*-namenode-*.log多半是重复格式化导致 clusterID 不一致少了 DataNode 多半是目录权限或磁盘空间问题。2.3 数据入湖从本地 CSV 到 HDFS环境起来后先把酒店 CSV 传上去。别用put直接怼大文件先建目录再传# 建分层目录 hdfs dfs -mkdir -p /warehouse/ods/hotel_raw hdfs dfs -mkdir -p /warehouse/dwd/hotel_detail hdfs dfs -mkdir -p /warehouse/dws/hotel_agg # 上传本地文件 hdfs dfs -put ./hotel_data.csv /warehouse/ods/hotel_raw/ # 验证 hdfs dfs -ls /warehouse/ods/hotel_raw/ hdfs dfs -tail /warehouse/ods/hotel_raw/hotel_data.csv-tail看最后 1KB用来确认文件没传断。如果 CSV 有表头后面建 Hive 表时用tblproperties(skip.header.line.count1)跳过别手动删表头删了原始层就不「原始」了。3. 用 Hive 把酒店数据清洗成可分析明细表3.1 建 ODS 外部表并核对字段Hive 建外部表指向 ODS 目录删表不删数据安全。假设 CSV 字段顺序是酒店ID、名称、省份、城市、区县、星级、评分、评论数、最低价、地址。CREATE EXTERNAL TABLE ods_hotel_raw ( hotel_id STRING, hotel_name STRING, province STRING, city STRING, district STRING, star STRING, score STRING, comment_cnt STRING, low_price STRING, address STRING ) ROW FORMAT DELIMITED FIELDS TERMINATED BY , STORED AS TEXTFILE LOCATION /warehouse/ods/hotel_raw TBLPROPERTIES (skip.header.line.count1);字段全用 STRING 是故意的——原始数据里数字列经常混着空值、单位、异常字符先当字符串接进来清洗阶段再转类型避免加载时因类型转换失败丢行。建完先SELECT * LIMIT 10看一眼确认分隔符和列数对得上。列数对不上是最常见的地址里带逗号会把一列劈成两列这时候要么换分隔符比如\t要么在采集端做转义。3.2 清洗规则空值、单位、口径统一酒店数据脏点集中在几处评分空或为「暂无」、价格带「」和「起」、省份带「省/市/自治区」后缀不统一、星级有「五星级」「5」「五」多种写法。清洗用 Hive SQL 一把梭CREATE TABLE dwd_hotel_detail STORED AS PARQUET AS SELECT hotel_id, trim(hotel_name) AS hotel_name, -- 省份口径统一去掉后缀 regexp_replace(trim(province), (省|市|自治区|特别行政区)$, ) AS province, trim(city) AS city, trim(district) AS district, -- 星级归一中文/数字统一成 1-5 整数 CASE WHEN star RLIKE [1-5] THEN CAST(regexp_extract(star, [1-5], 0) AS INT) WHEN star LIKE %五% THEN 5 WHEN star LIKE %四% THEN 4 WHEN star LIKE %三% THEN 3 WHEN star LIKE %二% THEN 2 WHEN star LIKE %一% THEN 1 ELSE NULL END AS star_level, -- 评分非数字置空 CASE WHEN score RLIKE ^[0-9](\\.[0-9])?$ THEN CAST(score AS DOUBLE) ELSE NULL END AS score, CAST(comment_cnt AS BIGINT) AS comment_cnt, -- 价格抽数字 CAST(regexp_extract(low_price, ([0-9]), 1) AS INT) AS low_price, trim(address) AS address FROM ods_hotel_raw WHERE hotel_id IS NOT NULL AND trim(hotel_id) ;逻辑说明regexp_replace统一省份口径让「广东省」和「广东」归并星级用CASE WHEN兜住多种写法匹配不到就置 NULL 而不是瞎猜价格用regexp_extract抽第一段数字把「300起」变成 300评分用正则卡住纯数字格式非数字置 NULL。WHERE过滤掉主键为空的行这是去重和保证后续关联正确的前提。参数上注意regexp_extract的第三个参数是捕获组序号写 0 是整段匹配写 1 是第一个括号组价格这里只有一个括号所以写 1。转 Parquet 时如果报内存不足调set hive.exec.dynamic.partition.modenonstrict;和适当提高mapreduce.map.memory.mb。3.3 去重与异常值处理同一家酒店可能重复采集按hotel_id去重保留最新一条假设有采集时间字段没有就保留任意一条CREATE TABLE dwd_hotel_dedup STORED AS PARQUET AS SELECT * FROM ( SELECT *, ROW_NUMBER() OVER (PARTITION BY hotel_id ORDER BY comment_cnt DESC) AS rn FROM dwd_hotel_detail ) t WHERE rn 1;用ROW_NUMBER而不是DISTINCT因为要按规则挑一条而不是简单去重。ORDER BY comment_cnt DESC是经验做法——评论数多的那条通常信息更全。异常值方面评分限定在 0-5价格限定在合理区间比如 0 到 100000超出范围的置 NULLSELECT COUNT(*) AS total, SUM(CASE WHEN score 0 OR score 5 THEN 1 ELSE 0 END) AS bad_score, SUM(CASE WHEN low_price 0 OR low_price 100000 THEN 1 ELSE 0 END) AS bad_price FROM dwd_hotel_dedup;先统计再决定阈值别拍脑袋。跑完这一步DWD 层就是干净可用的明细表了。4. 各省市酒店指标聚合从明细到能看的结论4.1 按省、按市的酒店数量与评分聚合DWS 层做聚合产出直接能进报表的指标。按省统计酒店数、平均评分、平均价格CREATE TABLE dws_hotel_by_province STORED AS PARQUET AS SELECT province, COUNT(DISTINCT hotel_id) AS hotel_cnt, ROUND(AVG(score), 2) AS avg_score, ROUND(AVG(low_price), 0) AS avg_price, SUM(comment_cnt) AS total_comments FROM dwd_hotel_dedup WHERE province IS NOT NULL GROUP BY province;按城市统计同理把GROUP BY换成city但要注意城市名跨省可能重名比如「城区」这种严谨做法是GROUP BY province, city。COUNT(DISTINCT hotel_id)比COUNT(*)稳防止上游没去干净。ROUND保留两位让结果好看也避免浮点尾数干扰对比。星级分布用一次透视SELECT province, SUM(CASE WHEN star_level 5 THEN 1 ELSE 0 END) AS star5, SUM(CASE WHEN star_level 4 THEN 1 ELSE 0 END) AS star4, SUM(CASE WHEN star_level 3 THEN 1 ELSE 0 END) AS star3, SUM(CASE WHEN star_level 2 THEN 1 ELSE 0 END) AS star_low FROM dwd_hotel_dedup GROUP BY province;这种写法比多次GROUP BY star_level再 join 更省一次扫描酒店数据字段不多时优势不明显但数据量上来后差别可观。4.2 用 MapReduce 做一遍词频式统计理解底层Hive 跑得爽但课程设计常要求手写 MapReduce。酒店数据里「各城市酒店数量」本质就是词频统计的变体。Mapper 输出city, 1Reducer 累加// Mapper切分 CSV取城市字段作为 key public class HotelCityMapper extends MapperLongWritable, Text, Text, IntWritable { private final Text city new Text(); private final IntWritable one new IntWritable(1); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line value.toString(); // 跳过表头 if (line.startsWith(hotel_id)) return; String[] fields line.split(,, -1); if (fields.length 4) return; // 列数不足直接丢 city.set(fields[3].trim()); // 城市在第 4 列 context.write(city, one); } }// Reducer累加同城计数 public class HotelCityReducer extends ReducerText, IntWritable, Text, IntWritable { private final IntWritable result new IntWritable(); Override protected void reduce(Text key, IterableIntWritable values, Context context) throws IOException, InterruptedException { int sum 0; for (IntWritable v : values) sum v.get(); result.set(sum); context.write(key, result); } }参数说明line.split(,, -1)的-1保证末尾空字段不被丢弃否则列数判断会误伤fields.length 4是防御性过滤脏行直接跳过而不是抛异常中断整个作业。Driver 里设job.setJarByClass、setMapperClass、setReducerClass输出类型Text/IntWritable要和 Mapper、Reducer 一致不一致会在运行时抛类型不匹配。打包提交hadoop jar hotel-city.jar com.demo.HotelCityDriver \ /warehouse/dwd/hotel_detail /output/hotel_city_cnt hdfs dfs -cat /output/hotel_city_cnt/part-r-00000 | head输出目录必须不存在否则作业直接失败这是 MapReduce 的保护机制重跑前先hdfs dfs -rm -r /output/hotel_city_cnt。4.3 结果导出与可视化衔接聚合结果从 HDFS 拉回本地接 Python 做可视化hdfs dfs -getmerge /warehouse/dws/hotel_by_province ./hotel_by_province.csv-getmerge把目录下多个 part 文件合并成一个本地文件省得手动拼。拿到 CSV 后用 pandas 读进来画图import pandas as pd import matplotlib.pyplot as plt df pd.read_csv(hotel_by_province.csv, headerNone, names[province, hotel_cnt, avg_score, avg_price, total_comments]) df df.sort_values(hotel_cnt, ascendingFalse).head(15) plt.figure(figsize(12, 6)) plt.barh(df[province], df[hotel_cnt]) plt.xlabel(hotel count) plt.title(Top 15 provinces by hotel count) plt.gca().invert_yaxis() plt.tight_layout() plt.savefig(hotel_by_province.png, dpi150)参数说明headerNone是因为 Hive 导出的 part 文件不带表头手动用names补head(15)只画前 15 名全画会挤成一团。这一步把 Hadoop 的批处理结果和 Python 可视化串起来课程设计里通常就是这条链路。5. 酒店数据跑 Hadoop 的避坑清单五个真实翻车现场5.1 现象作业卡在 map 100% reduce 0% 不动原因数据倾斜。某个省份或城市的酒店数量远超其他所有相同 key 落到同一个 Reducer单个 Reducer 扛不住。酒店数据里「北京」「上海」这种直辖市 key 的量级可能是小城市的几十倍。解决开启负载均衡set hive.groupby.skewindatatrue;Hive 会先做一轮随机分发预聚合或者在 key 上拼随机后缀打散聚合两次。MapReduce 里可以自定义 Partitioner 把大 key 拆开。5.2 现象中文省份名乱码聚合结果里出现问号原因CSV 编码是 GBKHive 默认按 UTF-8 读中文直接变乱码GROUP BY时同一个省被拆成多个乱码 key。解决上传前用iconv -f GBK -t UTF-8 hotel_data.csv hotel_data_utf8.csv转码或者建表时指定SERDEPROPERTIES(serialization.encodingGBK)。我一般统一转 UTF-8省得后面每个环节都要记编码。5.3 现象价格字段抽出来全是 NULL原因regexp_extract的正则没匹配上。原始价格可能是「300-500元」这种区间第一个数字是 300 没问题但如果写成「 300」和数字间有空格或者用了全角数字正则就抓不到。解决先SELECT low_price FROM ods_hotel_raw LIMIT 50肉眼看真实格式再写正则。全角数字要先translate转半角。别凭想象写正则这是血泪经验。5.4 现象伪分布式下 DataNode 起不来原因重复执行hdfs namenode -formatNameNode 的 clusterID 变了DataNode 目录里还是旧的 clusterID两者对不上直接拒绝启动。解决要么删掉 DataNode 数据目录重新格式化要么手动把dfs/data/current/VERSION里的 clusterID 改成和 NameNode 一致。根治办法是格式化前想清楚别反复 format。5.5 现象小文件太多后续查询慢原因MapReduce 或 Hive 动态分区写入时每个 task 产出一个文件分区多、task 多就产生大量小文件NameNode 内存压力和查询扫描开销都上去了。解决写入前set hive.merge.mapfilestrue;和set hive.merge.size.per.task256000000;让小文件合并已经产生的用ALTER TABLE ... CONCATENATEORC 格式或重跑一次合并作业。酒店数据按城市分区时特别容易踩这个城市几百个每个分区几个小文件。6. 把酒店分析做成可复跑的任务调度与增量技巧一次性跑完不算完真实场景数据会更新你得让这套流程能重复跑、能增量。我一般用 shell 脚本把「上传→清洗→聚合→导出」串起来再用 crontab 或调度工具定时触发。核心是每次跑之前清理目标目录、按日期分区写入避免数据叠加。#!/bin/bash set -e # 任何一步失败就退出别带着脏数据往下跑 DT$(date %Y%m%d) # 1. 上传当日数据到日期分区 hdfs dfs -mkdir -p /warehouse/ods/hotel_raw/dt${DT} hdfs dfs -put -f ./data/hotel_${DT}.csv /warehouse/ods/hotel_raw/dt${DT}/ # 2. 清洗写入 DWD 日期分区 hive -e SET hive.exec.dynamic.partitiontrue; SET hive.exec.dynamic.partition.modenonstrict; INSERT OVERWRITE TABLE dwd_hotel_detail PARTITION(dt${DT}) SELECT hotel_id, trim(hotel_name), regexp_replace(trim(province),(省|市)\$,), trim(city), trim(district), star, score, comment_cnt, low_price, address FROM ods_hotel_raw WHERE dt${DT}; # 3. 聚合 hive -e INSERT OVERWRITE TABLE dws_hotel_by_province PARTITION(dt${DT}) SELECT province, COUNT(DISTINCT hotel_id), ROUND(AVG(score),2) FROM dwd_hotel_detail WHERE dt${DT} GROUP BY province; # 4. 导出 hdfs dfs -getmerge /warehouse/dws/hotel_by_province/dt${DT} ./out/province_${DT}.csv echo done ${DT}set -e是关键缺了它某一步失败脚本还继续跑最后产出一份半截数据你还以为成功了。INSERT OVERWRITE配合分区重跑同一天不会重复累加天然幂等。分区字段dt放在SELECT最后Hive 动态分区要求分区列在查询结果末尾顺序错了直接报错。验证增量是否正确我习惯对比两天的行数和关键指标波动SELECT dt, COUNT(*) AS cnt, ROUND(AVG(avg_score),2) AS score FROM dws_hotel_by_province GROUP BY dt ORDER BY dt;如果某天行数突然翻倍多半是上游重复推送或分区没覆盖写如果平均评分断崖式变化回去查 DWD 层当天的清洗日志。这套「脚本 分区 幂等覆盖」的组合是我做离线数据处理最省心的习惯——不追求花哨的调度框架先把可复跑、可回溯做到位后面接什么调度工具都只是换个触发方式。数据这行能重跑出一样的结果比跑得快更让人踏实。希望帮到你。本文还有配套的精品资源点击获取
返回列表