简介本资源是一个基于Hadoop生态的实战型大数据分析项目面向Java开发者、大数据初学者及高校课程实践者聚焦海量酒店数据的分布式处理与统计分析。项目完整实现从HDFS数据存储到MapReduce编程的全流程涵盖全国各省市酒店信息如省份、城市、星级、价格等字段的清洗、分组聚合与多维统计助力掌握Hadoop核心组件与真实业务场景结合的关键能力。压缩包共79个文件含21个Java源码与21个编译后class文件构成可运行的MapReduce程序7个XML配置文件含Hadoop集群参数与Maven构建配置2个CSV数据源含hotel.csv原始数据以及说明.txt等关键文档整体体积仅758KB结构紧凑、开箱即用。已有2095人学习下载提供从本地测试到集群部署的完整代码工程、清晰目录组织含src/main/java、target/classes及hadoop-hotel脚本目录和可复现的分析逻辑是理解MapReduce编程范式与行业数据处理流程的优质入门范例。1. 用 Hadoop 处理全国酒店数据不是简单跑个 MapReduce而是构建可扩展、可复用、能支撑业务迭代的数据处理链路你手头有一份来自多个渠道的全国酒店数据——包含北京朝阳区某连锁酒店的房型价格、成都春熙路周边民宿的评分与评论文本、西安钟楼附近单体酒店的营业状态变更记录甚至还有部分地级市文旅局提供的备案信息 CSV。这些数据格式不一、更新频率不同、字段缺失严重但业务方明确要求按省份统计高星酒店数量、识别连续三个月评分下滑的酒店、生成各城市热门商圈的房价波动热力图。这时候Hadoop 不是“大材小用”的代名词而是唯一能承载这种多源异构、增量全量混合、需反复迭代清洗与特征计算的基础设施。它解决的不是“能不能算”而是“能不能在数据量从 GB 级涨到 TB 级时仍保持小时级产出、字段可追溯、逻辑可回滚”。适合正在从 Excel/MySQL 迁移至企业级数据平台的中台团队、需要对接 BI 工具做区域运营分析的酒店集团数据组以及承担政务文旅数据整合任务的信创项目组。本篇不讲 Hadoop 安装步骤只聚焦于如何让一份真实的全国酒店数据在 HDFS 上真正“活”起来——从原始文件落地到分层建模再到支撑下游分析任务。2. 基于 HDFS 的酒店数据分层存储设计用目录结构和命名规范替代“把所有 CSV 扔进一个文件夹”Hadoop 生态里数据组织方式直接决定后续开发效率与维护成本。面对全国酒店数据常见错误是把所有省市 CSV 直接hdfs dfs -put到/raw/hotel/下结果半年后没人能说清hotel_20231025_shanghai.csv和shanghai_hotel_v2.csv的差异。正确做法是建立严格分层目录体系并配合 Hive 外部表绑定路径让数据位置即语义。2.1 三层存储模型raw → clean → dim dwdraw 层原始层仅做最小化接入不做任何清洗或格式转换。目录结构按“数据来源 时间粒度 地域维度”组织/data/raw/hotel/gov_beijing/2024/06/28/ /data/raw/hotel/ota_meituan/2024/06/28/ /data/raw/hotel/ota_elong/2024/06/28/每个子目录下存放当日原始文件文件名含校验码如beijing_hotel_20240628_7a3f2d.csv.md5避免因网络中断导致文件不完整。clean 层清洗层对 raw 层数据执行字段标准化、空值填充、编码统一UTF-8、坐标系校准WGS84 → GCJ02。关键点在于清洗脚本必须输出与输入文件一一对应的 clean 文件且保留原始文件名前缀便于问题溯源# 清洗后路径示例 /data/clean/hotel/gov_beijing/2024/06/28/beijing_hotel_20240628_clean.csvdim dwd 层维度层 详细事实层这是面向分析的建模层。dim_province存储省级行政区划代码与名称映射dwd_hotel_info是核心宽表整合所有来源的酒店基础属性ID、名称、地址、经纬度、星级、开业时间、所属品牌并打上src_system来源系统标识和etl_batch_id本次 ETL 批次号两个关键字段为后续去重与变更识别提供依据。提示不要在 Hive 中用INSERT OVERWRITE TABLE dwd_hotel_info SELECT ... FROM raw_table直接覆盖全表。真实场景中酒店信息每日仅少量变更应采用INSERT INTO TABLE dwd_hotel_info PARTITION(dt20240628) SELECT ...按天分区插入并在下游任务中通过ROW_NUMBER() OVER (PARTITION BY hotel_id ORDER BY etl_batch_id DESC)取最新快照。这避免了全量重刷带来的资源浪费与窗口期数据不可用。2.2 Hive 表结构设计用分区分桶应对“省份查询高频、酒店 ID 查询低频”特性全国酒店数据最常被问的问题是“江苏有多少家四星以上酒店”、“杭州西湖区近三个月评分下降的酒店有哪些”。这意味着查询天然带有地域和时间过滤条件。Hive 表必须利用分区裁剪Partition Pruning能力CREATE EXTERNAL TABLE dwd_hotel_info ( hotel_id STRING, hotel_name STRING, province_code STRING, city_name STRING, district_name STRING, star_level TINYINT, score DOUBLE, lng DECIMAL(10,8), lat DECIMAL(9,8), brand STRING, src_system STRING, etl_batch_id STRING ) PARTITIONED BY (dt STRING) -- 按天分区如 20240628 CLUSTERED BY (province_code) INTO 32 BUCKETS -- 按省份分桶提升省份聚合性能 STORED AS PARQUET LOCATION /data/dwd/hotel/info/;PARTITIONED BY (dt STRING)使WHERE dt20240628查询自动跳过其他日期分区减少 I/O。CLUSTERED BY (province_code) INTO 32 BUCKETS将同一省份的酒店数据物理聚集在相同 HDFS 块内当执行GROUP BY province_code统计时Map 阶段本地性更高Shuffle 数据量显著降低。STORED AS PARQUET列式存储对SELECT COUNT(*) FROM dwd_hotel_info WHERE star_level 4 AND province_code 320000这类查询Parquet 只读取star_level和province_code两列比 TextFile 快 3–5 倍。2.3 数据接入自动化用 Oozie 调度 Shell 脚本完成“检测→清洗→加载”闭环手动hdfs dfs -put无法应对每日新增的数十个省市文件。需编写可配置的接入脚本并由 Oozie 定时触发#!/bin/bash # load_hotel_data.sh SOURCE_DIR/data/raw/hotel/$1/$2 # $1source_type, $2date TARGET_CLEAN/data/clean/hotel/$1/$2 TARGET_DWD/data/dwd/hotel/info # 1. 检查原始文件完整性MD5 校验 if ! md5sum -c $SOURCE_DIR/*.md5 --status; then echo MD5 check failed for $SOURCE_DIR 2 exit 1 fi # 2. 执行清洗调用 Spark SQL 脚本处理编码、空值、坐标 spark-sql \ --master yarn \ --conf spark.sql.adaptive.enabledtrue \ -f /opt/scripts/clean_hotel.sql \ --conf spark.sql.hive.convertMetastoreParquetfalse \ --conf spark.sql.files.ignoreMissingFilestrue \ -d INPUT_PATH$SOURCE_DIR \ -d OUTPUT_PATH$TARGET_CLEAN # 3. 加载至 DWD 层Hive INSERT INTO hive -e ALTER TABLE dwd_hotel_info ADD IF NOT EXISTS PARTITION (dt$2); INSERT INTO TABLE dwd_hotel_info PARTITION(dt$2) SELECT TRIM(hotel_id) as hotel_id, CASE WHEN LENGTH(TRIM(hotel_name))0 THEN UNKNOWN ELSE TRIM(hotel_name) END, get_province_code(address), -- UDF根据地址解析省码 city_name, district_name, CAST(COALESCE(star_level, 0) AS TINYINT), COALESCE(score, 0.0), CAST(lng AS DECIMAL(10,8)), CAST(lat AS DECIMAL(9,8)), COALESCE(brand, OTHER), $1 as src_system, $2 as etl_batch_id FROM clean_hotel_temp WHERE hotel_id IS NOT NULL AND LENGTH(TRIM(hotel_id)) 0; 该脚本被 Oozie 封装为 Coordinator每天凌晨 2 点扫描/data/raw/hotel/*/yyyy/MM/dd/目录动态生成 workflow 实例。关键参数$1来源类型和$2日期由 Oozie 的dateEL 函数传入实现“一份脚本多源复用”。3. 使用 Spark SQL 实现核心分析逻辑从“统计各省酒店数”到“识别评分异常酒店”的端到端代码HiveQL 适合简单聚合但酒店数据的深度分析如滑动窗口计算评分变化、基于地理围栏的商圈聚合必须依赖 Spark SQL 的 DataFrame API 与内置函数。以下代码均在spark-sqlCLI 或 Zeppelin 中可直接运行无需修改即可处理真实数据。3.1 按省份统计高星酒店数量用CASE WHENCOUNT避免多次扫描-- 计算各省四星及以上酒店数量含五星级 SELECT province_code, province_name, COUNT(*) AS high_star_count, ROUND(AVG(score), 2) AS avg_score FROM dwd_hotel_info a JOIN dim_province b ON a.province_code b.province_code WHERE a.dt 20240601 -- 限定最近一个月数据 AND a.star_level 4 GROUP BY province_code, province_name ORDER BY high_star_count DESC LIMIT 10;WHERE a.dt 20240601利用分区裁剪仅读取 2024 年 6 月 1 日后的数据避免扫描历史全量。JOIN dim_province维度表dim_province应设为TBLPROPERTIES(transactionaltrue)并启用 Hive ACID保证省份名称变更时能原子更新。ROUND(AVG(score), 2)对平均分保留两位小数符合业务报表精度要求。3.2 识别连续三个月评分下滑酒店用LAG()窗口函数实现跨月比较-- 步骤1先生成各省酒店月度平均分快照 CREATE OR REPLACE TEMPORARY VIEW hotel_monthly_score AS SELECT hotel_id, SUBSTR(dt, 1, 6) AS ym, -- 提取年月如 202406 ROUND(AVG(score), 2) AS monthly_avg_score FROM dwd_hotel_info WHERE dt 20240401 AND dt 20240630 -- 覆盖最近三个月 GROUP BY hotel_id, SUBSTR(dt, 1, 6); -- 步骤2用 LAG 计算前两个月分数标记连续下滑 SELECT hotel_id, ym, monthly_avg_score, prev1_score, prev2_score, CASE WHEN monthly_avg_score prev1_score AND prev1_score prev2_score THEN DOWN_3MONTHS ELSE NORMAL END AS trend_flag FROM ( SELECT hotel_id, ym, monthly_avg_score, LAG(monthly_avg_score, 1) OVER (PARTITION BY hotel_id ORDER BY ym) AS prev1_score, LAG(monthly_avg_score, 2) OVER (PARTITION BY hotel_id ORDER BY ym) AS prev2_score FROM hotel_monthly_score ) t WHERE ym 202406; -- 只看最新月份结果SUBSTR(dt, 1, 6)将分区字段dt如20240628转为202406作为月度粒度键。LAG(..., 1)获取同一酒店在上个月的平均分LAG(..., 2)获取上上个月。PARTITION BY hotel_id确保窗口按酒店分组。CASE WHEN ... THEN DOWN_3MONTHS业务规则直接嵌入 SQL输出结果可直接导入 BI 工具做告警。3.3 生成城市商圈房价热力图用ST_PointST_Contains实现地理围栏聚合假设已有一个dim_business_district表存储全国主要商圈的 WKT 多边形如POLYGON((116.4 39.9, 116.5 39.9, 116.5 39.8, 116.4 39.8, 116.4 39.9))-- 计算每个商圈内酒店的平均房价price 字段需提前从原始数据清洗出 SELECT b.district_name, b.city_name, COUNT(*) AS hotel_count, ROUND(AVG(a.price), 0) AS avg_price, MIN(a.price) AS min_price, MAX(a.price) AS max_price FROM dwd_hotel_info a JOIN dim_business_district b ON ST_Contains(ST_Polygon(b.wkt_polygon), ST_Point(a.lng, a.lat)) WHERE a.dt 20240628 AND a.price IS NOT NULL AND a.price BETWEEN 100 AND 20000 -- 过滤异常高价/低价 GROUP BY b.district_name, b.city_name ORDER BY avg_price DESC LIMIT 50;ST_Point(a.lng, a.lat)将经纬度转为 Spark GIS 的 Point 类型。ST_Polygon(b.wkt_polygon)将 WKT 字符串转为 Polygon。ST_Contains(...)判断酒店坐标是否落在商圈多边形内。这是空间 JOIN 的核心Spark 3.0 原生支持无需额外 GIS 库。a.price BETWEEN 100 AND 20000业务常识过滤排除录入错误如价格填成 0.01 元或 1000000 元。4. 处理酒店数据特有的脏数据问题编码混乱、地址歧义、星级标注不一致的实战方案全国酒店数据最大的技术挑战不在规模而在质量。某省文旅局 CSV 用 GBK 编码某 OTA 接口返回 UTF-8 但混入\x00控制字符某连锁酒店 Excel 中“五星”写成“★★★★★”这些都会导致 Hive 表加载失败或分析结果偏差。必须在清洗环节针对性解决。4.1 编码自动识别与转换用file命令 iconv预处理原始文件# 在清洗脚本开头加入编码探测 for file in /data/raw/hotel/*/2024/06/28/*.csv; do encoding$(file -i $file | sed s/.*charset\([^;]*\).*/\1/) if [[ $encoding ! utf-8 ]]; then iconv -f $encoding -t utf-8 $file -o ${file%.csv}_utf8.csv mv ${file%.csv}_utf8.csv $file fi donefile -iLinux 原生命令准确识别文件实际编码比单纯看文件后缀可靠。iconv标准编码转换工具支持 GBK、GB2312、BIG5 等中文常见编码。-f $encoding动态指定源编码避免硬编码导致转换失败。4.2 地址标准化用正则提取省市区再通过dim_province关联校验原始地址字段如广东省深圳市南山区科技园科苑路15号、上海静安区南京西路1266号恒隆广场需拆解为标准三级行政区划-- Hive UDFJava 实现get_province_code(address) public static String getProvinceCode(String address) { if (address null) return null; // 匹配省级前缀北京市、上海市、广东省、新疆维吾尔自治区... Pattern p Pattern.compile(^(北京市|上海市|天津市|重庆市|内蒙古自治区|广西壮族自治区|西藏自治区|宁夏回族自治区|新疆维吾尔自治区|香港特别行政区|澳门特别行政区|河北省|山西省|辽宁省|吉林省|黑龙江省|江苏省|浙江省|安徽省|福建省|江西省|山东省|河南省|湖北省|湖南省|广东省|海南省|四川省|贵州省|云南省|陕西省|甘肃省|青海省|台湾省)); Matcher m p.matcher(address); if (m.find()) { String provinceName m.group(1); // 查 dim_province 表返回 province_code如 110000 return lookupProvinceCode(provinceName); } return null; // 未匹配到交由人工复核 }正则表达式覆盖全部 34 个省级行政区全称包括“自治区”、“特别行政区”等后缀避免漏匹配。lookupProvinceCode()方法内部查询 Hive 的dim_province表确保代码与业务字典强一致而非硬编码映射。4.3 星级字段归一化用CASE WHENREGEXP_REPLACE统一为数字原始数据中星级表示法五花八门五星级★★★★★5星5*五星清洗 SQL 统一处理SELECT hotel_id, hotel_name, CASE WHEN star_text RLIKE ^[★*]{5}$ OR LOWER(star_text) RLIKE ^5[星*]$|^五星$ THEN 5 WHEN star_text RLIKE ^[★*]{4}$ OR LOWER(star_text) RLIKE ^4[星*]$|^四星$ THEN 4 WHEN star_text RLIKE ^[★*]{3}$ OR LOWER(star_text) RLIKE ^3[星*]$|^三星$ THEN 3 WHEN star_text RLIKE ^[★*]{2}$ OR LOWER(star_text) RLIKE ^2[星*]$|^二星$ THEN 2 WHEN star_text RLIKE ^[★*]{1}$ OR LOWER(star_text) RLIKE ^1[星*]$|^一星$ THEN 1 ELSE 0 -- 未评级或无效值 END AS star_level FROM raw_hotel_table;RLIKEHive 正则匹配比LIKE更灵活。LOWER(star_text)统一转小写避免大小写敏感问题。^5[星*]$匹配“5星”、“5*”^五星$匹配纯汉字覆盖主流写法。5. 性能调优与监控让 Hadoop 集群稳定支撑酒店数据日更任务当全国酒店数据日增 500 万行、单日清洗任务耗时从 12 分钟涨到 25 分钟时不能只靠加机器。需从数据、计算、资源三层面精准干预。5.1 数据层面用ANALYZE TABLE收集统计信息提升 Join 效率Hive 优化器依赖表的行数、列值分布等统计信息选择最优执行计划。对dwd_hotel_info执行ANALYZE TABLE dwd_hotel_info PARTITION(dt20240628) COMPUTE STATISTICS FOR COLUMNS hotel_id, province_code, star_level, score;COMPUTE STATISTICS FOR COLUMNS不仅收集表级行数还收集关键列的 NDV非重复值数量、最大最小值、空值率。后续JOIN dim_province ON a.province_code b.province_code时优化器知道province_code只有 34 个值会优先选择 Broadcast Join避免 Shuffle。5.2 计算层面Spark 参数调优针对酒店数据特点酒店数据宽表30 字段但单行不大 2KBShuffle 倾向于小文件多。关键参数参数推荐值说明spark.sql.adaptive.enabledtrue启用自适应查询执行AQE自动合并小 Task处理数据倾斜spark.sql.files.maxPartitionBytes128m控制每个 Partition 最大字节数避免单个 Task 处理过大文件spark.sql.autoBroadcastJoinThreshold50m维度表dim_province仅 100KB远小于此值强制 Broadcast Joinspark.serializerorg.apache.spark.serializer.KryoSerializerKryo 比 Java Serializer 快 2–3 倍序列化酒店对象更高效在spark-sql启动时传入spark-sql \ --master yarn \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.files.maxPartitionBytes134217728 \ --conf spark.sql.autoBroadcastJoinThreshold52428800 \ --conf spark.serializerorg.apache.spark.serializer.KryoSerializer \ -f analyze_hotel.sql5.3 监控层面用 YARN ResourceManager UI 定位慢任务根因当某个清洗任务卡在Stage 3: Shuffle Read时不要盲目重跑。登录 YARN Web UIhttp://rm-host:8088找到该 Application点击ApplicationMasterLogs搜索关键词GC overhead limit exceededJVM 内存不足需调大spark.executor.memoryFailed to connect to server网络问题检查 DataNode 是否存活Skipped 123456 records due to parse error原始文件存在非法字符回到file -i步骤重新探测编码。注意监控不是事后救火而是前置埋点。在清洗脚本末尾添加日志echo $(date %Y-%m-%d %H:%M:%S) INFO: Loaded $(hdfs dfs -count /data/clean/hotel/gov_beijing/2024/06/28 | awk {print $3}) files /var/log/hotel_etl.log这样每天可快速确认数据是否按时接入比等 BI 报表出错再排查快 6 小时。6. 用 Sqoop 实现酒店数据从 MySQL 到 HDFS 的增量同步避免全量导出的资源浪费很多酒店集团已有 MySQL 业务库其中hotel_info表每日新增/更新约 2 万条。若每次用sqoop import全量导出不仅占用 MySQL 从库 IO还会在 HDFS 上产生大量冗余文件。必须采用增量同步。6.1 基于--incremental append的主键增量模式前提MySQL 表有自增主键id且新数据id单调递增。sqoop import \ --connect jdbc:mysql://mysql-host:3306/hotel_db \ --username sqoop_user \ --password-file /user/sqoop/.pwd \ --table hotel_info \ --target-dir /data/raw/hotel/mysql_inc/20240628 \ --incremental append \ --check-column id \ --last-value 12345678 \ # 上次同步的最大 id --fields-terminated-by \001 \ --lines-terminated-by \n \ --null-string \\N \ --null-non-string \\N \ --compress \ --compression-codec org.apache.hadoop.io.compress.SnappyCodec \ --as-parquetfile--incremental appendSqoop 仅导出id 12345678的新记录。--as-parquetfile直接生成 Parquet 格式省去 Hive 后续INSERT INTO ... SELECT转换步骤。--compress --compression-codec ...使用 Snappy 压缩比默认 Gzip 速度快 3 倍压缩比略低但更适合 Hadoop 随机读。6.2 基于--incremental lastmodified的时间戳增量模式若 MySQL 表无可靠主键但有update_time字段类型为TIMESTAMP# 先查出上次同步的最新时间戳 LAST_TIME$(hive -e SELECT MAX(update_time) FROM dwd_hotel_info WHERE dt20240627 | tail -1) sqoop import \ --connect jdbc:mysql://mysql-host:3306/hotel_db \ --username sqoop_user \ --password-file /user/sqoop/.pwd \ --table hotel_info \ --target-dir /data/raw/hotel/mysql_time_inc/20240628 \ --incremental lastmodified \ --check-column update_time \ --last-value $LAST_TIME \ --merge-key id \ # 指定主键用于合并去重 --fields-terminated-by \001 \ --as-parquetfile--merge-key idSqoop 会将新数据与 HDFS 上旧数据按id合并自动覆盖update_time更新的记录实现“Upsert”语义。--last-value $LAST_TIMEShell 变量注入确保时间戳精确到秒。6.3 增量同步后的数据一致性校验用CHECKSUM验证 MySQL 与 Hive 数据行数同步完成后必须验证数据完整性-- 在 MySQL 中执行 SELECT COUNT(*) FROM hotel_info WHERE update_time 2024-06-28 00:00:00; -- 在 Hive 中执行 SELECT COUNT(*) FROM dwd_hotel_info WHERE dt 20240628;若两者相差超过 0.1%则触发告警。更严谨的做法是抽样校验关键字段-- Hive 中随机抽 100 条与 MySQL 对应 id 的记录比对 SELECT h.hotel_id, h.hotel_name, h.score, m.score AS mysql_score FROM ( SELECT hotel_id, hotel_name, score FROM dwd_hotel_info WHERE dt 20240628 DISTRIBUTE BY RAND() SORT BY RAND() LIMIT 100 ) h JOIN mysql_hotel_info m ON h.hotel_id m.id;只要发现h.score ! m.mysql_score就说明同步过程存在字段映射错误或时区转换问题需立即回滚并修复 Sqoop 参数。本文还有配套的精品资源点击获取