简介:基于Spark的外卖大数据平台分析系统完整项目包,主要面向大数据开发、数据分析和机器学习初学者,帮助解决外卖场景下的实时订单监控、用户行为分析与销量预测等问题。压缩包共含38个文件,以14个Scala源文件为核心,覆盖Spark Streaming、Spark SQL、MLlib等模块,并辅以5个Markdown说明文档、3个Hive SQL脚本、2个SQL脚本、2个JSON配置、1个Python脚本和Shell辅助脚本,还包含csv、tsv示例数据,整体大小仅646KB。项目完整呈现了从数据接入、清洗整合、特征工程到模型训练的典型流程,包含pom.xml与src/main标准工程目录,可帮助读者快速搭建外卖数据分析环境,并理解订单量预测、用户偏好挖掘、配送路径优化等实际应用场景。Scala源码中还包括参数配置与性能调优细节,便于深入学习Spark运行机制。目前已有740人学习下载,适合作为课程设计、毕业设计或入门实战的参考资料。 最近不少转行大数据的朋友都在找能练手、能写进简历的项目素材,“基于Spark的外卖大数据平台分析系统”这六个字在各类资源站里出现频率特别高。我拿了一个完整的项目压缩包,从解压开始一点点跑通,把离线统计和实时监控两条链路都搭了起来,再回头把Spark on YARN的CPU配额和内存模型调到一个相对可控的状态,这个过程里踩坑和反查的路径都很有代表性。这篇文章就把这套递进式的拆解思路整理出来:先是看系统和代码结构,然后是环境准备和导入工程,再是核心分析链路的实现逻辑,最后是集群资源调优和面试怎么讲。无论你是准备课程设计、数据开发岗位面试,还是已经开始接触公司里类似的大数据平台,里面的经验都能直接拿去用。
1. 拿到zip包后先看懂:外卖分析平台到底解决什么问题
1.1 外卖数据长什么样,为什么非要用Spark
外卖平台每天会产生大量订单记录,每一条都不是传统关系库里那种规整的行数据,而是一个嵌套很深的JSON。举个例子,一个普通外卖订单会同时包含用户ID、商家ID、骑手ID、城市ID、下单时间、支付金额、优惠金额、配送距离、配送时长、订单状态,以及一个菜品明细数组,数组里每个菜品又有SKU ID、名称、数量、单价。把这些放到一段JSON里,大概长这样:
{ "order_id": "202501150123456780001", "user_id": 1002345, "merchant_id": 88234, "rider_id": 500312, "city_id": 10, "order_time": "2025-01-15 12:03:11", "finish_time": "2025-01-15 12:47:02", "pay_amount": 36.8, "discount_amount": 8.2, "delivery_distance": 3.4, "delivery_duration": 32, "order_status": "finished", "items": [ {"sku_id": 1001, "name": "招牌黄焖鸡", "quantity": 1, "price": 26.5}, {"sku_id": 2033, "name": "米饭", "quantity": 2, "price": 3.2} ] }如果只有一个城市,一天可能就产生上百万条这样的数据,一个月下来几个亿条。此时再用MySQL去跑全局聚合,比如统计某区域一个月的GMV趋势、某个商家的用户复购率,查询性能会差到无法接受。这种场景天然适合Spark这类分布式计算框架:把几亿条数据切成多个分区,并行执行清洗、聚合、关联,最终把细粒度结果再回写到关系库供报表使用。
1.2 一个完整的外卖分析系统功能边界
市面上的“Spark外卖大数据分析系统”虽然项目名各异,但功能边界基本不出这三块:
- 离线数据看板:按日、周、月统计订单量、GMV、活跃用户数、活跃商家数,输出TOP N菜品、TOP N商家、TOP N区域,以及下单高峰时段分析。这是整个系统的地基,也是最容易出成果的部分。
- 实时监控链路:通过Kafka接入实时订单消息,用Spark Streaming做分钟级统计,比如最近5分钟下单量、当前实时销售额、超时配送告警。这块主要解决“看板数据落后一天”的问题。
- 用户画像与偏好推荐:基于用户历史订单,提取口味偏好、消费区间、常点商圈等标签,再产出菜品或商家推荐列表。做这块时要注意,推荐部分一般只是简单的规则统计,不会上太复杂的机器学习模型,但足以体现完整的数据处理能力。
把这三个模块拆清楚再看源码,就不会被各种文件搞晕。我在读代码时习惯先按“数据从哪来、经过什么处理、结果到哪去”这条链路去归类,项目里所有类基本都能归到这三个环节中。
2. 别急着解压,先把运行环境梳理一遍
2.1 版本选型和最小集群搭配
我见过不少人第一步就卡在环境上:JDK版本不对、Hadoop和Spark版本有兼容问题,或者干脆不知道这套东西要装哪些组件。先把一套稳妥的版本组合摆出来:
| 组件 | 推荐版本 | 说明 |
|---|---|---|
| JDK | 1.8 | Spark 3.x对JDK 8的支持最成熟,尽量不要用过于新的版本当学习环境 |
| Hadoop | 3.3.x | 对容器化、YARN资源隔离支持更好 |
| Spark | 3.2.x或3.3.x | 学习项目足够稳定,同时支持内置Hive |
| Kafka | 2.8.x | 配合Spark Streaming消费实时订单 |
| MySQL | 5.7或8.0 | 存放聚合结果、供大屏查询 |
| ZooKeeper | 3.7.x | Kafka和HDFS HA都可能依赖 |
如果只有一台笔记本,不用急着搭完整集群。先以local模式把代码跑通,把Spark默认的“本地跑”参数配上,数据换成小样本文件。等逻辑没问题了,再考虑搭一个一主两从的最小YARN集群。那种“先把集群搭起来再写代码”的顺序,对新手来说是效率最低的,因为环境问题会和代码问题搅在一起,出了问题很难定位。
2.2 zip解压和导入工程的报错处理
很多项目是zip压缩包分发,网络下载的压缩包损坏概率比你想象中高得多。我遇到过最典型的报错是:
导入资源包失败: caused by: invalid zip archive: could not find EOCDEOCD全称是End Of Central Directory Record,也就是ZIP格式尾部标识中央目录结束的那一小段数据。下载时文件被截断、网盘客户端同步中断、压缩包在Windows和macOS之间来回传输时编码出问题,都可能导致EOCD缺失。解决办法不是瞎试,而是先做完整性和结构测试:用7-Zip打开压缩包、执行“测试”功能,看能否完整列出目录;在Linux环境可以用unzip -t校验每个文件;有条件的话对比一下源文件的MD5和大小。一旦校验失败,最省事的就是重新下载,不要指望本地修复工具能救回来。
导入IDEA时也有一个很容易被忽略的细节:不要直接在Project面板里选择zip文件导入,而是先解压到本地目录,再以Maven或Gradle工程的方式打开。原因很简单,IDE需要拿到pom.xml或build.gradle才能识别项目结构,而很多错误的根因是Maven依赖下载不完整,本地仓库里的jar损坏了,此时IDE也会报“invalid zip archive”之类的错误。遇到这种情况,优先清理本地.m2/repository里的损坏目录,再重新reimport。
2.3 项目目录结构的正确阅读顺序
拿到了解压好的工程,别急着点开所有类。先看目录,一个典型的外卖分析平台目录大概长这样:
外卖大数据分析平台/ ├─ README.md ├─ docs/ # 设计文档、接口文档 ├─ sql/ # Hive建表、MySQL建表脚本 ├─ etl/ # 原始JSON清洗与转换 ├─ offline/ # 离线统计任务 ├─ realtime/ # Spark Streaming实时任务 └─ web/ # 可视化大屏页面阅读顺序建议是:README → sql目录 → etl目录 → offline目录 → realtime目录 → web目录。README会告诉你运行前提;sql目录让你知道目标表结构是怎么设计的,数据最终长成什么样;etl和offline是离线主链路;realtime负责实时补充。这个顺序符合数据流向,能让你在脑子里快速建立全貌。
3. 核心分析链路:从订单Json到可视化大屏
3.1 清洗与ETL:嵌套JSON扁平化
拿到手的原始订单JSON通常来自客户端埋点或服务端日志,里面会有各种脏数据:字段缺失、菜品数组为空、支付金额出现负值、订单状态为非法枚举值。所以第一步不是统计,而是清洗。
用Spark做清洗,核心是把JSON读成DataFrame,然后把嵌套结构展开成宽表。常见做法是先读成一行行文本,再用from_json加schema解析:
import org.apache.spark.sql.types._ val schema = new StructType() .add("order_id", StringType) .add("user_id", LongType) .add("merchant_id", LongType) .add("rider_id", LongType) .add("city_id", IntegerType) .add("order_time", StringType) .add("finish_time", StringType) .add("pay_amount", DoubleType) .add("discount_amount", DoubleType) .add("delivery_distance", DoubleType) .add("delivery_duration", IntegerType) .add("order_status", StringType) val raw = spark.read.text("hdfs:///data/orders/20250115/*.json") val orders = raw.select(from_json(col("value"), schema).alias("data")) .select("data.*") .filter(col("order_status").isin("finished", "canceled", "refunding")) .filter(col("pay_amount").gt(0))这里有两层考虑。第一是把JSON字段拍平,后续SQL不需要再嵌套解析;第二是提前过滤非法数据,避免统计出来的数字不可解释。我在实践中的体会是,清洗规则宁可先做保守的“保留合法数据”,也不要一上来就补默认值,否则业务方很容易质疑你的口径。
3.2 离线统计:Spark SQL典型写法
清洗完的数据进入离线统计层。这个环节最常用的就是Spark SQL,因为代码可读性强、也好配合Hive元数据。我自己的习惯是先把DataFrame注册成临时视图,再用SQL表达指标:
orders.createOrReplaceTempView("orders") val dailyStats = spark.sql( """ |SELECT date_format(order_time, 'yyyy-MM-dd') AS stat_date, | count(DISTINCT order_id) AS order_cnt, | round(sum(pay_amount), 2) AS gmv, | count(DISTINCT user_id) AS user_cnt, | count(DISTINCT merchant_id) AS merchant_cnt |FROM orders |WHERE order_status = 'finished' |GROUP BY date_format(order_time, 'yyyy-MM-dd') |""".stripMargin)如果需要做TOP N榜单,再加一层窗口函数,比如按城市统计销售额前三的商家:
SELECT city_id, merchant_id, gmv FROM ( SELECT city_id, merchant_id, round(sum(pay_amount), 2) AS gmv, row_number() OVER (PARTITION BY city_id ORDER BY sum(pay_amount) DESC) AS rn FROM orders WHERE order_status = 'finished' GROUP BY city_id, merchant_id ) t WHERE rn <= 3读代码的时候你会发现,离线模块很少写特别复杂的算子链,大量分析逻辑都是用SQL直译的。这正好说明一个理念:能用声明式表达清楚的分析,就不要用手写RDD算子去绕,Spark的Catalyst优化器会帮我们处理很多执行层面的细节。
3.3 实时链路:Spark Streaming消费Kafka
实时部分的核心是从Kafka消费订单消息,做微批次统计。项目里基于Spark Streaming的典型写法如下:
val kafkaParams = Map[String, Object]( "bootstrap.servers" -> "kafka1:9092,kafka2:9092", "key.deserializer" -> classOf[StringDeserializer], "value.deserializer" -> classOf[StringDeserializer], "group.id" -> "order_realtime_group", "auto.offset.reset" -> "latest", "enable.auto.commit" -> (false: java.lang.Boolean) ) val stream = KafkaUtils.createDirectStream[String, String]( ssc, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](Array("topic_orders"), kafkaParams) ) stream.map(_.value()).foreachRDD { rdd => if (!rdd.isEmpty()) { val df = spark.read.json(rdd) df.createOrReplaceTempView("orders_batch") val stat = spark.sql( """ |SELECT count(*) AS order_cnt, | round(sum(pay_amount), 2) AS gmv |FROM orders_batch |""".stripMargin) // 写入Redis/MySQL,用于大屏实时展示 } } ssc.start() ssc.awaitTermination()这里有一个很多初学者不理解的点:为什么Spark Streaming的窗口计算看起来没有“窗口”?因为foreachRDD本身就是以微批为单位在跑,每5秒或10秒一个批次,天然就是窗口。如果要统计“最近5分钟”的累计值,就需要自己维护一个状态,比如把每分钟聚合结果写入Redis,然后用Redis的过期时间或ZSET来做滑动窗口。这也是这类项目比纯离线更有工程感的地方。
3.4 结果落库与大屏展示
分析算完,最终要服务业务。离线结果和实时聚合结果会写入MySQL,典型的表结构包括:
| 表名 | 主要字段 | 用途 |
|---|---|---|
| ads_daily_stats | stat_date, order_cnt, gmv, user_cnt, merchant_cnt | 每日大盘概览 |
| ads_merchant_top | stat_date, city_id, merchant_id, gmv, rank | 商家排行 |
| ads_hour_trend | stat_date, hour, order_cnt, gmv | 高峰时段分析 |
| ads_realtime_order | window_start, order_cnt, gmv, avg_delivery_duration | 实时大屏数据 |
可视化层面,这类系统通常用Spring Boot搭一个轻量HTTP接口,前端用ECharts画大屏。核心图表包括GMV趋势折线、实时订单量滚动图、区域热力图、菜品排行条形图。注意,大屏只是一个结果展示层,真正的难点始终在上游:数据口径准不准、链路是否稳定、任务能不能在指定时间内跑完。
4. 集群资源问题实测:Spark on YARN的CPU与内存配置细节
4.1 Executor为什么会拿不到多核:问题复现
很多人在集群模式跑项目时发现,明明机器有16个CPU核,但Spark作业只用了1个核,任务慢到离谱,看YARN界面每个Executor只有一个vCore。这个现象的根因往往很直接:spark.executor.cores参数在YARN模式下的默认值是1。
Spark提交时如果不显式指定Executor核数,standalone模式下可能会尝试用满Worker可用核,但在YARN模式下走的是ApplicationMaster向ResourceManager申请容器,每个容器默认分配1个vCore。再加上很多入门集群没有修改yarn.nodemanager.resource.cpu-vcores,NodeManager报告给ResourceManager的核数并不是物理机真实核数,资源池本身就不够看。还有一个隐蔽因素是动态资源分配没开或者限得太小,导致Executor总数始终上不去。
4.2 内存模型不搞懂,OOM只是时间问题
CPU只分配1个核还能靠加参数解决,内存配置则需要更多理解。Spark在Executor内把堆内存分成若干区域,由统一内存管理器调度:spark.executor.memory是JVM堆大小;spark.memory.fraction决定堆内可用于执行和存储的比例,默认0.6;spark.memory.storageFraction是Storage在上面的占比,默认0.5。执行和存储之间可以互相借用,但执行被存储挤压时会主动淘汰缓存数据。
真正跑YARN模式时,容器总内存不是Executor堆大小,堆外还要算上spark.memory.offHeap.enabled对应的堆外内存,以及JVM自身开销、线程栈、元空间,这些统称为Overhead,默认约等于max(384MB, 0.1 * executor.memory)。如果不设置,Executor申请4GB时,实际容器内存可能是4.4GB甚至更多。我在排查OOM时发现,很多人只盯着executor.memory加大堆内存,忽视Overhead和堆外内存,结果容器直接超过YARN限制被杀。
常用内存参数整理如下:
| 参数 | 默认值 | 作用 |
|---|---|---|
| spark.executor.memory | 1g | 每个Executor的JVM堆内存 |
| spark.executor.memoryOverhead | max(384MB, 0.1*executor.memory) | 堆外开销内存 |
| spark.memory.fraction | 0.6 | 堆内可用于执行和存储的比例 |
| spark.memory.storageFraction | 0.5 | Storage区在fraction内的占比 |
| spark.memory.offHeap.size | 0 | 堆外内存大小,需配合开关 |
4.3 调参实战:从15分钟到5分钟
调参不要凭感觉,要有基准。我拿一份500万条订单的离线聚合任务做了对比,初始状态几乎全是默认值,跑完花了15分钟左右。之后调整了提交参数:
spark-submit \ --master yarn \ --deploy-mode cluster \ --class com.demo.offline.DailyStatsJob \ --driver-memory 2g \ --executor-memory 4g \ --executor-cores 4 \ --num-executors 10 \ --conf spark.memory.fraction=0.75 \ --conf spark.memory.storageFraction=0.4 \ --conf spark.dynamicAllocation.enabled=false \ --jars mysql-connector-java-8.0.26.jar \ offline-job.jar对比结果:
| 指标 | 调参前 | 调参后 |
|---|---|---|
| executor-cores | 1 | 4 |
| num-executors | 2 | 10 |
| shuffle分区控制 | 默认200 | 根据数据量设为400 |
| 任务总耗时 | 15分钟 | 5分钟左右 |
这个提升不是固定的,但能说明一个道理:很多Spark作业慢,不是代码写得差,而是资源配额压根没有匹配集群真实规格。建议你在做性能优化前,先写一个固定数据量的小程序做压测,记录基线,再改一个参数跑一次,这样才能知道每个参数的边际收益。
5. 项目吃透以后,面试官怎么问都不慌
5.1 把项目讲成有数据、有链路、有取舍
面试时,讲这个项目要避免“我做了个外卖分析平台,用了Spark”这种一句话概括。比较好的表达结构是:数据规模有多大、数据从哪里来、离线链路怎么处理、实时链路怎么处理、最终结果怎么用、遇到什么问题。我这里提供一个可以套用的骨架:
这个项目模拟外卖平台的核心数据流,通过Kafka接入线上订单消息,统一解析成标准订单宽表。离线部分基于Spark SQL做日、周、月维度的大盘统计和TOP榜单,数据落MySQL供大屏查询。实时部分用Spark Streaming做分钟级订单量和GMV统计。整个过程中我重点解决了数据倾斜和资源参数调优两个问题,也把YARN内存模型对容器大小的影响理清了。
这段描述听起来有厚度,因为它同时覆盖了数据源、计算引擎、存储层和问题维度。
5.2 高频追问与应对思路
围绕这个项目,面试官常见追问基本集中在以下问题:
| 问题 | 应对要点 |
|---|---|
| Spark宽窄依赖区别 | 窄依赖父分区只对应一个子分区;宽依赖涉及Shuffle,父分区对应多个子分区,典型如groupByKey |
| 数据倾斜怎么解决 | 两阶段聚合加盐、广播小表、拆分热点Key、调整并行度 |
| Spark Streaming和Structured Streaming区别 | 前者基于RDD微批、API较老;后者基于DataFrame、有内置窗口和水位线支持 |
| 为什么不用Flink | 学习项目技术栈统一,Spark流批一体,吞吐高;生产低延迟场景Flink更合适,关键是说明取舍意识 |
| checkpoint机制 | 把作业状态和元数据持久化到HDFS,恢复重启时从中断处继续 |
| 如何保证数据不丢 | Kafka offset手动管理、开启Spark Streaming WAL、结果写入幂等处理 |
这些追问背后真正考的不是概念本身,而是你有没有深入想过“出问题怎么办”。所以吃透项目最好的方式,不是背答案,而是真正在集群上制造几次故障,比如把并行度调低模拟数据倾斜,把内存改小观察OOM,再对着日志一步步排查。这个过程比项目本身更能帮你建立信心。
本文还有配套的精品资源,点击获取