基于Spark的外卖大数据平台分析系统:从数仓分层到实时看板
2026/9/16 11:58:10 网站建设 项目流程

简介:基于Spark的外卖大数据平台分析系统毕业设计项目,覆盖数据采集、清洗、存储、分析与可视化展示等典型环节,面向计算机相关专业毕业生、课程设计学生及Spark初学者,既可作为毕设/课设的完整参考,也可作为大数据技术实践入门项目。压缩包共41个文件,大小仅649KB,核心是14个Scala源码文件,同时包括6个Markdown文档、3个HSQL脚本、SQL数据库脚本、XML/JSON配置文件以及Python辅助工具脚本,结构清晰,方便按模块查阅。详细文档涵盖需求分析、系统设计、环境部署与使用说明,附带效果截图,便于快速理解项目整体架构与实现思路。目前已有350人浏览/学习,项目代码已在Windows和macOS环境下运行验证,并获导师认可、答辩评分95分,适合直接基于此项目扩展功能,也可为同类外卖或电商数据分析毕业设计提供重要参考。

1. 从外卖订单到经营决策:这套系统到底在解决什么问题

深夜十点,外卖平台的后台还在被订单写入轰炸:商家要看午高峰取消率,运营要按商圈看热力图,决策层想对比这个月和上个月的复购变化。如果直接拿 MySQL 跑这些聚合,上百万条订单明细会把多维关联拖到几十秒,同时报表口径经常对不上。把“订单明细清洗、分层、离线聚合、可视化”这条链路交给 Spark,就是“基于 Spark 的外卖大数据平台分析系统”这个题目的真正内容。它不是冷门选题,而是一套标准的数据仓库加离线计算工程模型,可复现阶段非常清晰:数据采集、HDFS 存储、Spark SQL 清洗、分层聚合、导出 MySQL、大屏展示。适合正在选毕业设计题目的数据专业学生,也适合想快速搭一套可演示大数据链路的开发新人——把项目当入门级生产方案看,比当应付答辩的作业看收获大得多。

2. 为什么选 Spark 做外卖分析:组件定位与数仓分层设计

拿到题目先别急着写代码,前三周的时间应该花在技术选型和数据分层上。外卖场景的数据特征,决定了你选用哪套技术栈:日订单量百万级、单条数据几十个字段、计算以 T+1 的聚合报表为主、偶尔要按小时看趋势。“百万到千万级日增、小时级延迟可接受”,这个区间里 Spark 是最稳的选择。

2.1 Spark 在外卖数据管道里的位置:和 MapReduce、Flink 的边界

很多毕业设计会把 Hadoop、Spark、Flink 混着写,答辩时被问到“为什么不用另外一个”就卡住了。我的建议是先把边界说清楚。

对比维度MapReduceSpark(本题目核心)Flink
计算模型磁盘落地的 map/reduce 两阶段DAG 内存计算流批一体
批任务延迟分钟到小时级秒到分钟级秒级(批)/ 毫秒级(流)
开发效率需要手写大量 MR 逻辑SQL 覆盖 80% 以上场景流处理 API 学习成本高
适合本题目吗能跑但效率低非常适合偏重,杀鸡用牛刀

外卖数据分析以离线批处理为主,Spark 的核心优势在三点:一是 DataFrame/DataSet API 做了大量算子下推优化,同样的 join 和 groupBy 不需要你手动控制 map 过程;二是 Spark SQL 天然支持 Hive 表,元数据可以直接复用,数仓分层逻辑不用自己重写一套;三是社区资料多,出问题能搜到明确答案。Flink 在实时维度上有优势,但如果只做 T+1 报表,引入流引擎只是增加集群负担。

2.2 能跑起来的集群拓扑:三台机器怎么分配角色

毕设环境不需要大集群,三台虚拟机或者三台云主机就够。常规做法是严格区分主节点和计算节点:主节点只跑 NameNode、ResourceManager 和 Spark 的调度角色,计算节点才分配 Executor。这样避免一个进程吃掉所有资源导致任务排队。

节点角色分配推荐配置
masterNameNode + ResourceManager + Spark Master(可选)8C 16G
slave1DataNode + NodeManager + Executor4C 8G
slave2DataNode + NodeManager + Executor4C 8G

这里不要忽略数据本地性:Executor 优先调度到存有该数据块的 DataNode 上。如果 master 上也放了 Executor,跨节点拉数据的概率会变大,小数据量看不出来,跑到几千万行 join 时网络开销就明显了。生产环境建议把 master 的 yarn.nodemanager.resource.memory-mb 调低或直接不跑计算容器。

2.3 数仓四层怎么承接外卖业务:ODS、DWD、DWS、ADS

题目叫“平台分析系统”,底子是数仓。没有分层的 Spark 项目,最终一定会变成“一张大宽表”加上一堆互相覆盖的临时任务。哪怕只做毕业设计,也按 ODS、DWD、DWS、ADS 四层建库。

  • ODS 层:订单、商家、用户、骑手等原始 JSON/日志,原样落地。
  • DWD 层:数据清洗、维度退化,把城市名、商圈名直接灌进事实表。
  • DWS 层:按业务主题做轻度聚合,比如“商家-小时”“城市-天”“用户-签约月份”。
  • ADS 层:面向大屏和报表的最终结果集,通常已经是几十行的宽表。

建表语句用 Hive 语法写在 Spark SQL 里执行:

CREATE EXTERNAL TABLE ods_order_info ( order_id STRING, user_id STRING, merchant_id STRING, city_id INT, order_amount DECIMAL(10,2), delivery_fee DECIMAL(10,2), order_status STRING, order_ts BIGINT, ts TIMESTAMP ) PARTITIONED BY (dt STRING) STORED AS PARQUET LOCATION 'hdfs://master:8020/data/ods/order_info';

订单时间戳字段同时保留 order_ts(原始毫秒时间戳)和 ts(解析后的可读时间),是因为后续按小时聚合时,直接对 ts 做 hour() 比在 SQL 里反复 from_unixtime 更高效。选择 Parquet 而不是 ORC 的一个实际理由:Spark 对 Parquet 的谓词下推和列剪裁支持更成熟,且 Python 生态(比如后续要读回 Pandas 画图)兼容性更好。分区字段 dt 必须显式指定。

3. ETL 与指标聚合:用 Spark SQL 把订单明细变成经营指标

选型和分层定完后,进入核心开发:把原始数据加工成指标。这一章我按一条完整的跑批链路来讲,从读取、清洗到聚合,最后提交到 YARN 上。新手最容易犯的错,是上来直接写一个大脚本做全流程,结果中间某个字段格式错了,整个任务失败重跑半小时。

3.1 第一道工序:把 JSON 日志读进 DataFrame 并落成 Parquet

外卖系统的原始数据一般是埋点日志,按天落地到 HDFS,文件格式经常是压缩的 JSON。第一步先把它读进来,做一次纯粹的格式转换。

from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("etl_order_json") \ .config("spark.sql.adaptive.enabled", "true") \ .enableHiveSupport() \ .getOrCreate() df_raw = spark.read.json("hdfs://master:8020/raw/order/2025-01-0*.json") df_raw.write \ .mode("overwrite") \ .partitionBy("dt") \ .format("parquet") \ .saveAsTable("ods.order_info") print(f"load times: {df_raw.count()}")

注意读取时用 glob 表达式把一天内多个小时的文件一次扫进来,Spark 会并行拉取。count() 在这里不只是打印行数,它的执行会触发一次 Action,从而验证 Schema 推断是否成功;如果 JSON 里混入脏数据导致类型推断失败,这一步会直接暴露。写入 ODS 时使用 partitionBy("dt"),dt 字段必须提前存在于 DataFrame 中——这也是为什么原始埋点必须带上业务日期。

3.2 清洗与维度退化:DWD 层在 Spark 里做什么

DWD 的清洗要点不是“去重”这么简单。外卖订单里有几个典型脏场景:订单金额为负的测试订单、订单状态为 UNKNOWN 的中间态、merchant_id 为空的日志。还要做维度退化,把 city 名称、商圈名称直接 join 进来,避免后续每层都跟维表反复关联。

from pyspark.sql.functions import col, when, hour, from_unixtime df_ods = spark.table("ods.order_info") df_ods_clean = df_ods.filter( (col("order_amount") > 0) & (col("merchant_id").isNotNull()) & (col("order_status").isin("FINISHED", "CANCELLED", "REFUNDED")) ).withColumn( "order_hour", hour(from_unixtime(col("order_ts") / 1000)) ).withColumn( "is_cancel", when(col("order_status") == "CANCELLED", 1).otherwise(0) ) df_ods_clean.write \ .mode("overwrite") \ .partitionBy("dt") \ .format("parquet") \ .saveAsTable("dwd_order_info")

过滤条件里保留 CANCELLED 不是因为它是有效订单,而是为了后续算“取消率”这类比 GMV 更有业务价值的指标。order_hour 在 DWD 层就计算好,因为 DWS 聚合时按小时分桶是高频操作,把这个字段提前物化能省掉每次聚合都做的 from_unixtime 转换。

3.3 建立 DWS 聚合模型:商家 x 时段 x 城市:Spark SQL 的 groupBy 怎么写

DWS 层是“预聚合”的核心层。外卖平台的高频指标是:每个商家每天每个小时的订单量、GMV、取消量、客单价;每个城市每天的回落情况。这里用 DataFrame API 聚合一次,结果写回 DWS 表。

df_dws_merchant_hour = df_ods_clean.groupBy( "merchant_id", "city_id", "dt", "order_hour" ).agg( count("*").alias("order_cnt"), sum("order_amount").alias("gmv"), sum("is_cancel").alias("cancel_cnt"), round(avg("delivery_fee"), 2).alias("avg_delivery_fee") ) df_dws_merchant_hour.write \ .mode("overwrite") \ .partitionBy("dt") \ .saveAsTable("dws.merchant_hour_stat")

groupBy 之后 Spark 一定会触发 shuffle,这正是美团这类头部平台订单量下性能瓶颈的来源之一。这里 groupBy 的粒度是“商家 x 城市 x 小时”,因为大屏展示的“热力图”和“高峰趋势”都需要这个粒度;如果只做 T+1 总报表,groupBy 到天就够了,会快很多。取舍原则:DWS 粒度越粗,下游查询越快,但灵活度越低。对于外卖场景,“小时”是最小可接受粒度,再细到分钟,数据量会大几倍且无实际业务意义。

3.4 提交到集群:spark-submit 的参数怎么给

本地 IDE 里跑通了不代表集群能跑通,毕业设计的高频翻车点全在提交参数上。我一般用 spark-submit 配合 YARN cluster 模式:

spark-submit \ --master yarn \ --deploy-mode cluster \ --name order-etl-daily \ --driver-memory 2g \ --executor-memory 4g \ --executor-cores 2 \ --num-executors 4 \ --conf spark.sql.shuffle.partitions=200 \ --conf spark.sql.adaptive.coalescePartitions.enabled=true \ --py-files /home/bi/deps.zip \ hdfs://master:8020/app/jobs/order_etl.py \ --date 2025-01-08
参数作用说明
--deploy-mode clusterDriver 运行在 YARN 的 ApplicationMaster 内客户端提交后即可断开,适合定时任务
--num-executors 4Executor 数量与 NodeManager 可用核数保持匹配
--executor-memory 4g单个 Executor 堆内存别超过容器物理内存的 80%
--conf spark.sql.shuffle.partitionsShuffle 分区数默认 200,小集群反而要调小

注意 --py-files 传的依赖包路径,在 cluster 模式下本地文件不会被 YARN 识别,必须使用 HDFS 路径或提前分发到所有节点。任务代码也同理,hdfs:// 路径而不是 /home/bi/order_etl.py。

4. 调优与排错:Spark 跑外卖批任务最常踩的 6 个点

集群任务跑通了只是第一步。外卖数据有个特点:订单量集中在头部商家,高峰集中在午晚两餐——这种“长尾分布”和“时段聚集”直接决定了你调优的方向。有几个点必须在答辩前亲手踩一遍,能讲清楚就是加分项。

4.1 shuffle 是外卖订单关联的命门:先试 broadcast join

订单明细表 join 商家维表是最高频操作。按常规做法要先问:维表多大?如果商家表只有几万行,几百 MB 级别,直接用 broadcast join,把维表复制到每个 Executor 的内存里,省掉整个 shuffle。

from pyspark.sql.functions import broadcast df_merchant = spark.table("dim.merchant") # 约几万行,小于 10MB df_order_merged = df_ods_clean.join( broadcast(df_merchant), on="merchant_id", how="left" )

Spark 的 autoBroadcastJoinThreshold 默认是 10MB,超过这个阈值得手动指定。调优时先看执行计划是 SortMergeJoin 还是 BroadcastHashJoin,如果显示 SortMergeJoin 且维表确实小,那就是阈值没配好。对 5 年以上经验的人提醒一句:broadcast 是把维表全量复制到每个 Executor,如果维表 2GB 且有 50 个 Executor,总内存开销是 100GB,不是 2GB。

4.2 数据倾斜:头部商家占了一半订单

外卖订单天然倾斜:一个商圈头部商家的日订单量可能是尾部商家的几百倍。groupBy merchant_id 时,某个 reduce 任务挂着跑不完,其他任务都 idle 了。加盐是标准解法。

from pyspark.sql.functions import rand, floor, explode, array, lit # 大表侧:给订单加随机盐 df_big_salted = df_ods_clean \ .withColumn("salt", (floor(rand() * 10)).cast("int")) # 维表侧:把商家维表膨胀 10 份 df_dim_salted = df_merchant \ .withColumn("salt", explode(array([lit(i) for i in range(10)]))) # 关联时两个 key 一起 join df_joined = df_big_salted.join( df_dim_salted, ["merchant_id", "salt"], "left" ).drop("salt")

加盐后同一个商家的订单被打散到 10 个 Task 里并行处理,每个 Task 只处理原 task 的十分之一数据。代价是维表膨胀了 10 倍,所以盐的份数要根据倾斜程度调整:倾斜不严重时加 4 份就好,加太多反而会让 broadcast 成本盖过收益。

4.3 内存参数配错,任务一半就失败

小集群最常见的内存翻车是:给 Executor 设了 8g 内存,但 NodeManager 只有 8g 物理内存,yarn 容器直接起不来。容器实际占用是 executor-memory 加 spark.memory.overhead(默认 384MB)。有一个经验值:executor-memory 上限不超过机器物理内存除以该机器最大并发 Executor 数,再打八折。

参数默认值本集群建议值说明
spark.executor.memory1g4g单个 Executor 堆内存
spark.memory.fraction0.60.6(保持)执行与存储共享内存占比
spark.memory.storageFraction0.50.5Storage 内存占比,缓存 RDD 用
spark.sql.shuffle.partitions20064分区数 = Executor 数 x cores x 4 左右
spark.serializerJava 序列化KryoSerializer减少内存占用,但需注册类

4.4 spark on yarn 提交是不是只需要一个 spark 客户端:区别与条件

这道题在面试中也高频出现。答案是:提交端确实只需要一个装了 Spark 客户端程序的机器,它负责把 Application 提交给 YARN,但不一定参与计算。区别在于集群侧必须满足两点:第一,HDFS 里要有 Spark 依赖的 jar,或者 NodeManager 节点本身有 Spark 安装包;第二,--deploy-mode client 时 Driver 跑在客户端 JVM 里,而 cluster 时跑在 YARN 的 ApplicationMaster 容器里。

对比项client 模式cluster 模式
Driver 位置客户端机器YARN 容器内
适合场景交互式调试、本地跑小任务定时调度、生产环境
日志查看直接打印在终端需要 yarn logs -applicationId 拉取

毕设的每日跑批调度建议用 cluster 模式,因为如果用 client 模式,调度平台或者 cron 退出了 Driver 也一起挂。

5. 从结果到交付:大屏数据、近实时扩展与任务健康检查

5.1 聚合结果导出 MySQL,供数据可视化大屏接口查询

DWS/ADS 层的聚合结果一般不会直接给前端,而是落进 MySQL,后端提供一个只读接口,大屏前端再通过接口拉数据。Spark 写 MySQL 的常规做法是用 JDBC 写入。

df_dws = spark.table("dws.merchant_hour_stat") df_dws.write \ .mode("overwrite") \ .format("jdbc") \ .option("url", "jdbc:mysql://master:3306/meituan_bi") \ .option("dbtable", "t_merchant_hour_stat") \ .option("user", "bi_user") \ .option("password", "******") \ .option("batchsize", "2000") \ .option("truncate", "true") \ .save()

每一次跑批后先清空目标表再写入,避免脏数据累积。注意 dbtable 不要直接写全量库表名,可以写成子查询格式(select * from t where dt='2025-01-08') t,这样只覆盖当天的分区数据。batchsize 控制 JDBC 批量写入大小,2000 是比较稳的经验值,过大会占满 MySQL 连接。

5.2 想升级成近实时看板:Structured Streaming 接入的边界

如果导师或评委问“能不能看实时单量”,就用 Structured Streaming 接 Kafka。这里要说清:它不是从离线数仓改出来,而是另起一套实时链路,和 Spark 批任务并行跑。

df_stream = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "master:9092") \ .option("subscribe", "order-topic") \ .load() df_order = df_stream.select( from_json(col("value").cast("string"), order_schema).alias("o") ).select("o.*") df_gmv = df_order.groupBy( window(col("ts"), "1 minute"), col("city_id") ).agg(sum("amount").alias("gmv")) df_gmv.writeStream \ .outputMode("append") \ .format("console") \ .option("truncate", "false") \ .start() \ .awaitTermination()

这里有一个高频误区:Structured Streaming 的 groupBy 会做增量状态管理,状态会无限增长,所以 window 一定得带上,并且配置 state 清理策略。对于外卖业务,“1 分钟窗口”是最合理的粒度,再细到秒级,Kafka 的数据量和状态膨胀成本都会翻倍。

5.3 每日批任务的健康检查:退出码、分区与对账脚本

批任务最怕不是跑挂,而是“跑成功了但数据是错的”。用退出码判断会漏掉脏数据问题,所以要在调度脚本里加对账逻辑。下面是一个 health check 的 shell 骨架:

#!/bin/bash dt=$1 spark-submit --master yarn --deploy-mode cluster --date $dt exit_code=$? if [ $exit_code -ne 0 ]; then echo "Spark job failed, exit_code=$exit_code" exit 1 fi # 1. 检查分区存在 partition_cnt=$( hive -e "show partitions dws.merchant_hour_stat" | grep "dt=$dt" | wc -l ) # 2. 对账:当天行数与昨天行数比较 today_cnt=$(spark-sql -e "select count(*) from dws.merchant_hour_stat where dt='$dt'") yesterday_cnt=$(spark-sql -e "select count(*) from dws.merchant_hour_stat where dt='$prev_dt'") diff_ratio=$(echo "scale=4; ($today_cnt - $yesterday_cnt) / $yesterday_cnt" | bc) if [ -z "$partition_cnt" ] || \ [ "$today_cnt" -eq 0 ] || \ [ $(echo "$diff_ratio > 0.2" | bc) -eq 1 ]; then echo "Data check failed: today_cnt=$today_cnt, yesterday_cnt=$yesterday_cnt" exit 2 fi echo "Health check passed."

把退出码、分区数、行数差三个维度上报到监控系统或简单的 webhook,对比单看 exit_code 可靠得多。商家量级、节假日活动、数据源断流的场景下,目标行数都可能跳变,但对账脚本能通过“相对差异”而不是“绝对数量”把异常暴露出来。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询