简介:本资源是一个基于Hadoop生态的美团外卖大数据分析实战项目,面向大数据初学者与高校课程实践者,聚焦真实业务场景下的分布式数据处理能力训练。项目完整复现了用户行为、餐厅运营、物流配送等多维度分析流程,涵盖HDFS存储、MapReduce计算、HiveQL查询及自定义Partitioner等核心技能点。压缩包共89个文件,含48个Java主程序(如ProvincePartitionDriver、ReduceSideJoin等)、9个XML配置文件、7个CSV样本数据(如meituan.csv、us-counties.csv)、7个可执行JAR包及Shell脚本test.sh,辅以HTML报告页与可视化图片,整体7.37MB,结构清晰、模块解耦,便于分步调试与功能验证。已有91人学习下载,读者可直接运行MR任务、复现TopN统计、侧连接分析、序列文件写入等典型操作,并通过源码理解Hadoop组件协同机制与业务逻辑映射关系。
1. 基于Hadoop的美团外卖数据分析:不是跑个WordCount就叫大数据,而是真实还原订单洪峰、骑手调度瓶颈与区域供需失衡的离线分析闭环
你下载的这个基于Hadoop的美团外卖数据分析.zip,不是教学Demo,也不是PPT式“大数据演示”,它是一套可直接在伪分布式Hadoop环境上跑通的、带完整业务逻辑链路的离线分析工程。它用真实脱敏的美团外卖结构化日志(含订单时间、商户ID、用户ID、骑手ID、配送距离、超时标记、菜品品类、地理围栏编码),完成了从HDFS原始数据摄入 → MapReduce清洗去重 → Hive建模分层(ODS/DWD/DWS)→ Spark SQL聚合统计(日单量趋势、热力区域TOP20、骑手接单响应时长分布、品类复购率漏斗)→ 最终导出CSV供BI工具接入的全链路。适合两类人:一是正在准备大数据开发岗面试、需要一个能讲清“为什么用MapReduce而不用Spark做清洗”“Hive分区怎么设才不OOM”的项目背书者;二是刚搭好Hadoop伪分布式环境、苦于找不到有业务语义的真实数据练手的工程师。它不依赖ZooKeeper高可用集群,但所有脚本和配置都预留了YARN资源队列、HDFS权限控制、Hive事务表等企业级扩展点——换句话说,你今天在单机上跑通,明天就能无缝迁到三节点集群。
2. 数据架构设计与Hadoop组件选型:为什么用MapReduce清洗+Hive建模+Spark聚合,而不是全Spark或全Flink?
2.1 业务数据特征决定技术栈组合:高吞吐写入、低频更新、强SQL分析需求
美团外卖日志具备典型离线分析场景特征:
- 写入密集、读取稀疏:订单日志按分钟级批量写入HDFS,单日可达GB级,但分析任务每天只跑1–2次;
- 更新极少、删除无:订单状态变更通过追加新记录实现(如“已下单→配送中→已完成”),无需实时更新;
- 分析维度固定、SQL友好:业务方最常问的是“朝阳区上周奶茶类订单超时率”“骑手平均接单响应时间TOP10商圈”,天然适配Hive/Spark SQL;
- 清洗逻辑复杂、需强一致性:去重需跨天合并(同一订单可能因网络重试产生多条日志),且必须保证“同一订单ID只保留最早一条有效记录”,MapReduce的Shuffle阶段天然支持全局Key排序,比Spark的
distinct()更可控(后者在数据倾斜时易OOM)。
提示:这不是技术教条,而是血泪经验——我曾用Spark Streaming实时处理类似日志,结果因上游Kafka消息重复、下游MySQL幂等写入失败,导致报表数据每日偏差3%。离线批处理反而更稳。
2.2 组件分工明确:MapReduce负责“脏数据手术刀”,Hive负责“数据仓库骨架”,Spark负责“分析加速器”
| 组件 | 承担角色 | 关键配置/参数说明 | 为何不可替代 |
|---|---|---|---|
| MapReduce | 原始日志清洗:去重、字段校验(时间戳格式、金额正数)、异常订单过滤(配送距离>50km)、生成唯一订单ID | mapreduce.map.memory.mb=2048,mapreduce.reduce.memory.mb=4096;Reducer数设为min(20, 总输入文件数)防小文件 | HDFS小文件合并、跨Partition全局排序能力远超Spark RDD,尤其适合“按订单ID分组取最早时间戳”这类操作 |
| Hive | 构建分层模型:ODS层(原始日志,按dt=20240501分区)、DWD层(清洗后事实表+维度表关联)、DWS层(宽表预聚合,如dws_order_daily_by_area) | 启用ORC格式+ZLIB压缩(STORED AS ORC TBLPROPERTIES("orc.compress"="ZLIB")),分区字段dt类型为STRING(避免Hive 3.x对DATE类型分区的兼容问题) | Hive Metastore提供统一元数据管理,业务方用Beeline直连即可查,无需写代码;且Hive on Tez比MapReduce快3–5倍,适合DWS层轻量聚合 |
| Spark SQL | 高性能聚合分析:计算区域热力图(经纬度转GeoHash)、骑手响应时长分布(直方图bin=30s)、品类复购率(窗口函数ROW_NUMBER() OVER(PARTITION BY user_id ORDER BY order_time)) | spark.sql.adaptive.enabled=true,spark.sql.adaptive.coalescePartitions.enabled=true;Driver内存设为4g,Executor内存8g | Spark Catalyst优化器对复杂JOIN+GROUP BY的执行计划优于Hive,且DataFrame API比HiveQL更易调试(支持.explain()看物理计划) |
2.3 文件结构与目录规范:HDFS路径设计直接影响后续ETL可维护性
解压ZIP后,你会看到如下HDFS部署结构(需手动创建):
# 必须提前创建的HDFS目录(用hdfs dfs -mkdir -p) /user/hive/warehouse/ # Hive默认warehouse路径 /data/meituan_raw/ # 原始日志存放处,按dt分区 /data/meituan_clean/ # MapReduce清洗后输出路径 /data/meituan_dws/ # DWS层宽表输出路径 /tmp/spark_output/ # Spark临时输出(会自动清理)关键细节:
- 原始日志命名规则:
meituan_order_20240501_000001.log.gz(日期+序号),确保MapReduce InputFormat能正确识别Gzip压缩; - Hive外部表指向:
CREATE EXTERNAL TABLE ods_meituan_order (...) LOCATION '/data/meituan_raw/';—— 外部表避免误删HDFS数据; - DWS层分区策略:
dws_order_daily_by_area表按dt STRING, area_code STRING双分区,area_code为5位GeoHash编码(如wx4g7),查询“朝阳区”时只需扫描对应分区,避免全表扫描。
3. 核心脚本详解与执行流程:从HDFS数据上传到最终报表生成的六步实操
3.1 第一步:准备Hadoop伪分布式环境(验证HDFS/YARN/Hive可用)
注意:本项目要求Hadoop 3.3.6 + Hive 3.1.3 + Spark 3.3.2,版本不匹配会导致JAR包冲突。ZIP包内
env_check.sh已封装检测逻辑:
#!/bin/bash # env_check.sh:检查Hadoop/Hive/Spark基础服务 hdfs dfs -ls / >/dev/null 2>&1 && echo "✅ HDFS OK" || echo "❌ HDFS not running" yarn node -list 2>/dev/null | grep RUNNING >/dev/null && echo "✅ YARN OK" || echo "❌ YARN not running" hive -e "show databases;" >/dev/null 2>&1 && echo "✅ Hive OK" || echo "❌ Hive not running" spark-sql --version 2>/dev/null | grep "3.3.2" >/dev/null && echo "✅ Spark OK" || echo "❌ Spark version mismatch"执行后若全部✅,继续;否则请先完成《Hadoop伪分布式搭建全过程》(ZIP包内附hadoop_setup_guide.pdf,含Ubuntu/Windows WSL双平台图文步骤)。
3.2 第二步:上传原始数据到HDFS并验证分区
原始日志位于ZIP包/data/raw/目录下,共7天数据(20240501至20240507),每份约120MB:
# 创建HDFS原始数据目录并上传(逐日上传,避免单次传输超时) for dt in {20240501..20240507}; do hdfs dfs -mkdir -p /data/meituan_raw/dt=$dt hdfs dfs -put ./data/raw/meituan_order_${dt}_*.log.gz /data/meituan_raw/dt=$dt/ done # 验证上传完整性:检查每个分区文件数与大小 hdfs dfs -ls /data/meituan_raw/dt=20240501 | wc -l # 应返回12(12个Gzip文件) hdfs dfs -du -h /data/meituan_raw/dt=20240501 | awk '{sum += $1} END {print sum}' # 应≈120M逻辑说明:-put命令会自动将本地文件上传至HDFS,dt=$dt是Hive分区约定写法;-du -h用于校验数据量,避免网络中断导致部分文件未传全。
3.3 第三步:运行MapReduce清洗Job(核心:去重+字段标准化)
清洗逻辑封装在/src/main/java/com/meituan/etl/CleanOrderMapper.java中,关键逻辑:
// CleanOrderMapper.java 片段:提取关键字段并打标 public void map(Object key, Text value, Context context) throws IOException, InterruptedException { String line = value.toString(); String[] fields = line.split("\\|"); // 美团日志以|分隔 if (fields.length < 8) return; // 字段数不足跳过 String orderId = fields[0].trim(); String orderTime = fields[1].trim(); // 格式:2024-05-01 08:32:15 String distance = fields[5].trim(); // 配送距离(km) // 强校验:订单时间必须为合法timestamp,距离必须为数字且<50 if (!isValidTimestamp(orderTime) || !isNumeric(distance) || Double.parseDouble(distance) > 50.0) { context.getCounter("CLEAN", "INVALID_RECORD").increment(1); return; } // 输出:key=orderId, value=orderTime|distance|...(用于Reducer去重) context.write(new Text(orderId), new Text(orderTime + "|" + distance + "|" + ...)); }提交Job命令(在/script/目录下):
# clean_mr.sh:提交清洗Job hadoop jar target/meituan-etl-1.0.jar \ com.meituan.etl.CleanOrderDriver \ -D mapreduce.job.name="Meituan_Clean" \ -D mapreduce.map.memory.mb=2048 \ -D mapreduce.reduce.memory.mb=4096 \ -files /script/clean_config.xml \ /data/meituan_raw/ \ /data/meituan_clean/参数说明:
-D设置JVM内存,防止OOM;-files指定配置文件(含正则表达式规则),避免硬编码;- 输入路径
/data/meituan_raw/为HDFS路径,Job会自动遍历所有dt=子目录; - 输出路径
/data/meituan_clean/为清洗后数据,注意:此路径必须为空,否则Job失败。
3.4 第四步:Hive建模——创建三层表并加载数据
在/script/hive_ddl.sql中执行建表语句(Beeline连接):
-- 1. ODS层:原始日志(外部表) CREATE EXTERNAL TABLE IF NOT EXISTS ods_meituan_order ( order_id STRING, order_time STRING, user_id STRING, shop_id STRING, rider_id STRING, distance DOUBLE, is_timeout BOOLEAN, category STRING ) PARTITIONED BY (dt STRING) ROW FORMAT DELIMITED FIELDS TERMINATED BY '|' STORED AS TEXTFILE LOCATION '/data/meituan_raw/'; -- 2. DWD层:清洗后事实表(内部表,ORC格式) CREATE TABLE IF NOT EXISTS dwd_meituan_order ( order_id STRING, order_time TIMESTAMP, user_id STRING, shop_id STRING, rider_id STRING, distance DOUBLE, is_timeout BOOLEAN, category STRING, geo_hash STRING -- 由经纬度计算得出,用于区域聚合 ) PARTITIONED BY (dt STRING) STORED AS ORC TBLPROPERTIES("orc.compress"="ZLIB"); -- 3. 加载数据:从清洗后路径导入DWD层(按分区加载) INSERT OVERWRITE TABLE dwd_meituan_order PARTITION(dt='20240501') SELECT order_id, CAST(order_time AS TIMESTAMP), user_id, shop_id, rider_id, distance, is_timeout, category, geo_hash_from_latlon(lat, lon) as geo_hash -- UDF函数,ZIP包内已编译 FROM ( SELECT *, SUBSTR(order_time, 1, 10) as dt -- 从order_time提取日期 FROM ods_meituan_order WHERE dt='20240501' ) t;关键点:
geo_hash_from_latlon()是自定义UDF(Java实现),已打包进hive-udf.jar,需先ADD JAR hive-udf.jar;INSERT OVERWRITE会覆盖分区,确保数据一致性;CAST(order_time AS TIMESTAMP)将字符串转为Hive TIMESTAMP类型,支持后续时间窗口函数。
3.5 第五步:Spark SQL聚合分析(生成业务指标)
/script/spark_analyze.py使用PySpark执行:
from pyspark.sql import SparkSession from pyspark.sql.functions import * spark = SparkSession.builder \ .appName("Meituan_Analysis") \ .config("spark.sql.adaptive.enabled", "true") \ .getOrCreate() # 读取DWD层数据(按日期范围) df = spark.read \ .option("mergeSchema", "true") \ .orc("hdfs://localhost:9000/data/meituan_dwd/") \ .filter(col("dt").between("20240501", "20240507")) # 计算区域热力图:按geo_hash统计订单量 hot_area_df = df.groupBy("geo_hash") \ .agg(count("*").alias("order_cnt")) \ .orderBy(desc("order_cnt")) \ .limit(20) # 导出为CSV(注意:生产环境应写入HDFS,此处为本地调试) hot_area_df.coalesce(1).write.mode("overwrite").csv("./output/hot_area_top20.csv")执行命令:
spark-submit \ --master yarn \ --deploy-mode client \ --driver-memory 4g \ --executor-memory 8g \ --num-executors 4 \ --conf spark.sql.adaptive.enabled=true \ /script/spark_analyze.py逻辑说明:
coalesce(1)强制合并为1个文件,避免生成part-00000等碎片;--master yarn表示提交到YARN集群,若本地测试可改local[*];- 输出CSV位于本地
./output/,可直接用Excel打开。
3.6 第六步:验证结果与业务解读
最终生成的hot_area_top20.csv内容示例:
geo_hash,order_cnt wx4g7,12485 wx4g8,11932 wx4g6,10876 ...对照GeoHash编码表(ZIP包内/doc/geo_hash_mapping.csv),wx4g7对应北京市朝阳区CBD核心区。结合rider_response_time_distribution.csv(骑手接单响应时长分布),发现该区域0-30s响应占比仅62%,低于全市均值78%——印证了“订单密集但骑手运力不足”的业务假设。这就是离线分析的价值:不是罗列数字,而是定位根因。
4. 避坑指南:六个真实翻车现场与血泪解决方案
4.1 现象:MapReduce Job卡在ACCEPTED状态,YARN Web UI显示Application Status为ACCEPTED但不启动
原因:YARN资源不足,yarn.scheduler.maximum-allocation-mb默认值(8192MB)小于Job申请的mapreduce.reduce.memory.mb=4096,但NodeManager实际可用内存被其他进程占用,导致Container无法分配。
解决:
- 查看NodeManager日志:
tail -100 $HADOOP_HOME/logs/yarn-*-nodemanager-*.log | grep -i "memory"; - 降低Job内存申请:
-D mapreduce.reduce.memory.mb=2048; - 或调大YARN配置:在
yarn-site.xml中增加
<property> <name>yarn.scheduler.maximum-allocation-mb</name> <value>16384</value> </property>并重启YARN。
4.2 现象:Hive建表成功,但SELECT COUNT(*) FROM dwd_meituan_order返回0,且DESCRIBE FORMATTED dwd_meituan_order显示Location: null
原因:建表时未指定LOCATION,Hive自动创建在默认warehouse路径(/user/hive/warehouse/),但数据实际写入/data/meituan_dwd/,元数据与物理路径不一致。
解决:
- 方案A(推荐):建表时显式指定
LOCATION '/data/meituan_dwd/'; - 方案B:用
ALTER TABLE dwd_meituan_order SET LOCATION '/data/meituan_dwd/';修正元数据; - 切记:修改LOCATION后需执行
MSCK REPAIR TABLE dwd_meituan_order;同步分区。
4.3 现象:Spark SQL执行GROUP BY geo_hash时出现java.lang.OutOfMemoryError: Java heap space
原因:GeoHash基数过大(全国编码超百万),groupBy触发Shuffle时单个Reducer内存溢出。
解决:
- 启用AQE(Adaptive Query Execution):
spark.sql.adaptive.enabled=true(已默认开启); - 增加Shuffle分区数:
spark.sql.adaptive.coalescePartitions.enabled=false+spark.sql.files.maxPartitionBytes=128m; - 或改用近似算法:
df.approxQuantile("geo_hash", Array(0.5), 0.01)替代精确统计。
4.4 现象:clean_mr.sh执行报错ClassNotFoundException: com.meituan.etl.CleanOrderDriver
原因:Hadoop classpath未包含编译后的JAR包,或JAR包未打包依赖(如commons-lang3)。
解决:
- 检查JAR包是否含依赖:
jar -tf target/meituan-etl-1.0.jar | grep lang; - 若无,则用Maven Shade Plugin重新打包:
<plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <version>3.4.1</version> <executions> <execution> <phase>package</phase> <goals><goal>shade</goal></goals> </execution> </executions> </plugin>4.5 现象:Hive查询dwd_meituan_order时提示Failed to load data from path 'hdfs://...' due to permission denied
原因:HDFS目录权限为drwx------(仅owner可读),而Hive服务运行用户(如hive)无权限访问/data/meituan_dwd/。
解决:
- 修改HDFS权限:
hdfs dfs -chmod -R 755 /data/meituan_dwd/; - 或修改Owner:
hdfs dfs -chown -R hive:hive /data/meituan_dwd/; - 安全提示:生产环境应使用HDFS ACL而非宽松权限。
4.6 现象:Spark导出CSV时生成多个part-*.csv文件,且文件名含随机UUID,无法直接用Excel打开
原因:Spark默认按分区并行写入,coalesce(1)未生效(因DataFrame未触发Action)。
解决:
- 确保
coalesce(1)后紧跟write:df.coalesce(1).write.csv("path"); - 或改用
repartition(1)(但会引发全量Shuffle,慎用); - 终极方案:用
spark.sql("SELECT * FROM table").coalesce(1).write.mode('overwrite').option('header','true').csv('path')。
5. 进阶技巧:如何用这套框架快速适配你的业务数据?三个可复用的改造模板
5.1 模板一:替换数据源——从美团日志切换到你自己的订单系统
你不需要重写整个Pipeline,只需修改三处:
- 数据格式适配:修改
CleanOrderMapper.java中的split("\\|")为你的分隔符(如,或\t),并调整fields[]索引; - 字段映射配置:在
/conf/column_mapping.json中定义你的字段名到标准字段的映射:
{ "your_order_id": "order_id", "your_create_time": "order_time", "your_distance_km": "distance", "your_is_timeout_flag": "is_timeout" }- UDF扩展:若你的数据含经纬度,复用
geo_hash_from_latlon();若含用户画像标签,新增UDFget_user_segment(),编译后ADD JAR即可。
我从那以后每次接手新数据源,都先花1小时写这个JSON映射表,再跑通清洗Job——比重写MapReduce快10倍,且保证字段语义一致。
5.2 模板二:增加实时监控——用Hive Streaming对接Flink(非侵入式增强)
当前是纯离线,但业务需要“订单超时率超过15%时告警”。不必推翻重做,用Hive Streaming:
- 在DWD层表上启用Streaming:
ALTER TABLE dwd_meituan_order SET TBLPROPERTIES( "streaming.enabled"="true", "streaming.checkpoint.interval"="300000" );- 启动Flink Job监听Hive表变更(ZIP包内
/flink/stream_alert.py):
# 监听DWD表INSERT事件,计算5分钟滑动窗口超时率 table_env.execute_sql(""" CREATE TABLE hive_orders ( order_id STRING, order_time TIMESTAMP(3), is_timeout BOOLEAN, dt STRING ) WITH ( 'connector' = 'hive', 'table-name' = 'dwd_meituan_order', 'hive-conf-dir' = '/opt/hive/conf' ) """) # 后续接TUMBLING WINDOW聚合...这样,离线Pipeline不变,实时能力叠加其上。
5.3 模板三:性能调优速查表——针对不同硬件配置的参数组合
| 场景 | 推荐配置 | 依据 |
|---|---|---|
| 8核16G笔记本(伪分布式) | mapreduce.map.memory.mb=1024,mapreduce.reduce.memory.mb=2048,spark.executor.memory=4g,spark.executor.cores=2 | 避免Swap,留4G给系统 |
| 16核32G测试服务器(3节点) | yarn.nodemanager.resource.memory-mb=24576,mapreduce.reduce.memory.mb=6144,spark.executor.memory=12g,spark.executor.instances=3 | NodeManager内存设为总内存75%,Executor占NodeManager 50% |
| 生产集群(32核64G×10节点) | yarn.scheduler.capacity.root.default.maximum-capacity=80,hive.tez.container.size=8192,spark.sql.adaptive.enabled=true | 容量调度器限制default队列最大80%,Tez Container与YARN Container对齐 |
从那以后我每次部署新环境,都先查这张表,再
hadoop fs -du -h /看HDFS实际负载,最后微调——没再因OOM半夜被电话叫醒过。
希望帮到你。
本文还有配套的精品资源,点击获取