简介:基于SpringBoot与Spark的共享单车数据存储系统毕业设计项目,内含完整论文与可运行源码,适用于计算机科学与技术、人工智能等专业学生完成毕业设计或课程实践。包体共374个文件,以Java后端、Vue前端及SVG图标资源为主,并包含SQL初始化脚本、YML配置、Maven打包配置及一键启动批处理,压缩包大小17.8MB,目录层次清晰,便于按模块检索。系统围绕共享单车数据采集、存储与分析进行设计,借助Spark实现大规模数据查询、统计与趋势预测,通过SpringBoot构建RESTful接口,配套论文详细阐述了架构选型、数据库设计与实现方案,有助于读者贯通大数据处理与Web开发知识。已有50人学习该资源,下载后可根据README与论文快速还原项目,直接用于课题参考或二次开发。
1. 共享单车数据存储,为什么偏偏是 SpringBoot 加 Spark
做过几年后端的人应该都有体感:共享单车这类业务,数据量一旦跑起来,单机 MySQL 根本扛不住。每辆车每几秒上报一次位置,一天就是上亿条轨迹点;再加上订单、骑行时长、计费明细,日增数据轻松到 T 级。更麻烦的是,这些数据不是存下来就完事,还要支撑高峰期调度、潮汐分析、用户画像这些实时性要求不低的查询。传统关系型数据库在这种场景下,写入和查询都会被拖垮。
这套系统给出的解法是 SpringBoot 做业务层和应用接口,Spark 做分布式计算和批量分析,底层存储落在 HDFS 上。SpringBoot 的价值在于快速搭出 RESTful API,把前端、管理端、数据处理任务串起来;Spark 的价值在于面对海量骑行记录时,能用分布式内存计算把聚合分析从分钟级压到秒级。对做毕业设计或者课程项目的开发者来说,这个组合既覆盖了 Java Web 开发的完整链路,又接触到了大数据生态的真实工作方式,属于性价比很高的选题方向。
本文会沿着「数据链路设计 → Spark 分析模块落地 → SpringBoot 集成 → 部署优化 → 验证技巧」这条线展开,每一步都给出能直接用的代码和参数配置。
2. 数据链路设计:从单车传感器到 HDFS 分层存储
共享单车数据存储系统的第一个核心问题是:数据从哪里来,到哪里去,中间经过哪些处理。先把这个链路理清楚,后面的代码实现才不会跑偏。
2.1 数据采集端的多协议接入设计
共享单车的车锁终端通常通过 MQTT 或 HTTP 长连接上报数据,上报内容包含车辆 ID、经纬度、速度、电量、锁状态、时间戳等字段。这些数据的特点是频率高、字段固定、单条体积小,但总量巨大。
常见的做法是让终端先把数据推到消息队列,再异步写入存储层,避免终端请求直接压垮后端服务。系统在采集层用到 Kafka 作为缓冲区,SpringBoot 服务接收终端上报后,将原始 JSON 写入 Kafka topic,再由独立的消费者任务批量拉取并落盘。
// Kafka 生产者配置:批量发送,提高吞吐 @Bean public ProducerFactory<String, String> producerFactory() { Map<String, Object> props = new HashMap<>(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "node01:9092,node02:9092"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); // 批量发送,攒够 16KB 或等待 20ms 再发 props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); props.put(ProducerConfig.LINGER_MS_CONFIG, 20); props.put(ProducerConfig.ACKS_CONFIG, "1"); return new DefaultKafkaProducerFactory<>(props); }这段配置里有几个参数值得注意。BATCH_SIZE_CONFIG和LINGER_MS_CONFIG是配套使用的,前者控制攒多少字节才发一次,后者控制最多等多少毫秒。对共享单车这种高频小消息场景,默认的立即发送会导致网络包太多,吞吐上不去;调成批量模式后,单台 Broker 的写入性能能提升好几倍。ACKS_CONFIG设为1表示 Leader 写入即返回,兼顾了吞吐和可靠性,对位置数据这种允许极端情况下丢几条的场景足够。
2.2 存储层选型:HDFS 的目录规划与文件格式
数据落到 HDFS 后,目录规划直接决定后续 Spark 读取的效率。如果所有数据堆在一个目录下,Spark 做分区剪枝时会扫描大量无关文件,任务跑得又慢又费资源。
推荐按「业务类型 / 日期 / 小时」三层建目录,比如/bike/order/2025/06/01/和/bike/gps/2025/06/01/。日期和小时分区是共享单车数据分析里最常用的过滤维度——查高峰期、查某天的潮汐现象,都只需要扫描对应分区的文件。
文件格式上,不要直接存 JSON 或 CSV。JSON 没法做列剪枝和谓词下推,CSV 没有 schema 信息,Spark 读取时都要整文件扫一遍。生产环境常见的做法是存 Parquet,列式存储配合压缩,查询时只读需要的列,I/O 大幅减少。
-- Spark SQL 建表语句:Parquet 格式,按天分区 CREATE TABLE bike_order ( order_id STRING, bike_id STRING, user_id STRING, start_time TIMESTAMP, end_time TIMESTAMP, start_lng DOUBLE, start_lat DOUBLE, end_lng DOUBLE, end_lat DOUBLE, distance DOUBLE, amount DOUBLE ) USING parquet PARTITIONED BY (dt STRING);这份建表语句中,USING parquet让 Spark 把表数据以列式格式存储,查询时只反序列化涉及的列。PARTITIONED BY (dt STRING)是性能关键点,dt 分区字段既用于物理目录隔离,也是查询时下推条件的入口。
2.3 数据分层:原始层、明细层、聚合层
很多人在设计存储系统时,把原始数据和统计分析结果混在一起存,后面做报表发现要跑全量数据,越跑越慢。更规范的做法是把数据分成三层。
- 原始层:Kafka 里的 JSON 原样落盘,保留所有字段,用于回溯和排查问题
- 明细层:清洗、去重、格式转换后的 Parquet 文件,按天分区,是 Spark 分析的主数据源
- 聚合层:预计算好的统计结果,比如每小时的用车量、每个站点的周转率,直接供 SpringBoot 接口查询
聚合层的数据量小,可以存回 MySQL 或者直接用 Spark 的saveAsTable存成 Hive 表。SpringBoot 查询聚合层时响应时间在毫秒级,不需要每次实时跑 Spark 任务。
3. Spark 分析模块实战:骑行时长、潮汐规律与热点区域计算
存储系统搭好之后,真正的价值在分析。这一章用 Spark 的 DataFrame API 和 Spark SQL 完成几个共享单车场景里最常见的分析任务,每一段代码都给出运行逻辑和参数说明。
3.1 环境初始化与读取 HDFS 分区数据
Spark 读取分区数据时,一个常见的坑是直接读整个表的路径,导致全表扫描。正确做法是显式指定分区过滤条件,让 Spark 利用分区剪枝只读取需要的数据。
val spark = SparkSession.builder() .appName("BikeShareAnalysis") .master("yarn") .config("spark.sql.parquet.binaryAsString", "true") .config("spark.sql.shuffle.partitions", "200") .enableHiveSupport() .getOrCreate() // 只读取指定日期的订单数据 val orderDF = spark.read.parquet("/bike/order/2025/06/01")代码中spark.sql.shuffle.partitions是影响运行效率的关键参数,默认值是 200。如果集群规模小,shuffle 分区数太多会导致每个任务处理的数据量过小,调度开销反而变大;如果数据量大而分区数太少,单个任务会内存溢出。我一般按「总数据量 / 128MB」来估算,数据量在 10GB 左右时设为 100 到 200 之间。
3.2 单车使用峰值时段统计
运营方最关心的问题是:一天中哪些时段用车量最大,以便提前调度车辆。这个统计用 Spark SQL 的group by加hour函数就能完成,但要注意时间字段的类型转换。
orderDF.createOrReplaceTempView("orders") val peakHourDF = spark.sql(""" SELECT hour(start_time) AS hour_of_day, COUNT(*) AS order_count FROM orders WHERE dt = '2025-06-01' GROUP BY hour(start_time) ORDER BY order_count DESC """) peakHourDF.show(24)这段 SQL 中的hour(start_time)需要start_time是 Timestamp 类型,如果原始数据里是 String,必须先用to_timestamp做转换,否则会返回空结果。聚合结果只有 24 行,数据量极小,可以直接 collect 回驱动端,然后通过 SpringBoot 接口返回给前端图表。
3.3 热点区域识别:网格聚合与密度排序
热点区域分析是共享单车调度系统的基础。常见做法是把城市地图切成网格,统计每个网格内的租车和还车数量,找出供需失衡的区域。
// 把经纬度映射到 500m 左右的网格 val gridDF = orderDF .withColumn("grid_x", floor(col("start_lng") * 100)) .withColumn("grid_y", floor(col("start_lat") * 100)) .groupBy("grid_x", "grid_y") .agg( count("*").alias("rent_count"), sum("distance").alias("total_distance") ) .filter(col("rent_count") > 100) // 过滤掉随机零散订单floor(col("start_lng") * 100)这段是把经纬度放大 100 倍后向下取整,相当于把 0.01 度左右的范围归为一个网格,在中纬度地区大约对应 1 公里见方。网格大小可以按城市规模调整,城区密集区域用 0.005,郊区用 0.02。过滤条件rent_count > 100很关键,能去掉那些只是路过、没有调度价值的零散点。
这个任务的 shuffle 量比较大,groupBy("grid_x", "grid_y")会导致全量数据重分区。如果集群资源紧张,建议先按dt过滤到单天数据,或者用salting技术给热点 key 加随机前缀再分两次聚合。
3.4 用户骑行行为分析:RFM 模型简化版
用户行为分析需要把订单表按用户聚合,计算每个用户的骑行频次、平均时长和常用出发区域。这里的核心是避免数据倾斜——少数骑行达人可能有几千条记录,直接groupBy("user_id")会导致某个任务处理几百万行,其他任务空闲。
val userStatsDF = orderDF .groupBy("user_id") .agg( count("*").alias("ride_count"), avg(unix_timestamp(col("end_time")) - unix_timestamp(col("start_time"))).alias("avg_duration"), collect_set("grid_x").alias("frequent_grids") ) .filter(col("ride_count") >= 5)unix_timestamp计算骑行时长时要小心时区问题,HDFS 里存的时间如果带时区偏移,直接做差会有 8 小时的误差。建议在数据入仓时统一转成 UTC,展示层再转本地时间。collect_set会收集用户去过的所有网格,如果某个用户的网格数量非常大(超过几千),driver 端可能会内存压力大,可以用sort_array加limit限制只保留前几个常去区域。
4. SpringBoot 集成层实现:RESTful API 与 Spark 任务的协调
数据分析和业务接口之间需要一层稳定的桥梁。SpringBoot 在这一层负责三件事:暴露查询接口、触发定时分析任务、管理数据源的连接配置。
4.1 配置管理:分离业务库与统计结果库
共享单车系统的数据访问有个特点:业务数据在 Spark/HDFS 上,查询结果在 MySQL 里。如果所有数据源都写在application.yml里,容易混在一起。推荐的配置结构是分数据源管理,HDFS 的连接参数走 SparkConf,MySQL 只存聚合结果。
spring: datasource: dynamic: primary: business datasource: business: url: jdbc:mysql://node03:3306/bike_business username: bike_app password: ${BIKE_DB_PASSWORD} analytics: url: jdbc:mysql://node04:3306/bike_analytics username: bike_ana password: ${BIKE_ANA_PASSWORD}配置里把business和analytics拆成两个库,business存放车辆状态、用户账号等在线业务数据,analytics只存放 Spark 写入的统计结果。这样做的原因是两者的访问模式完全不同:业务库要求低延迟、强一致,统计库允许批量写入、偶尔读到旧数据。用@DS("analytics")注解切换到统计库即可。
4.2 查询层:返回聚合数据的接口实现
热点区域和峰值时段的结果在 Spark 任务跑完后写入analytics库,SpringBoot 接口直接查表返回。
@RestController @RequestMapping("/api/analytics") public class AnalyticsController { @Autowired private JdbcTemplate analyticsJdbcTemplate; @GetMapping("/peak-hours") public ApiResult<List<Map<String, Object>>> getPeakHours( @RequestParam String date) { String sql = "SELECT hour_of_day, order_count " + "FROM peak_hour_stats " + "WHERE dt = ? " + "ORDER BY order_count DESC"; List<Map<String, Object>> result = analyticsJdbcTemplate.queryForList(sql, date); return ApiResult.success(result); } }这里有几个实践要点。第一,SQL 里直接用?占位符,不拼接字符串,避免注入风险。第二,peak_hour_stats表在 Spark 写入时就应该按dt创建分区表,这样查询dt = '2025-06-01'时 MySQL 也能走分区裁剪。第三,如果前端需要实时查询某小时的数据,可以在 Redis 里缓存最近 24 小时的结果,过期时间设为 5 分钟。
4.3 任务触发:定时调用 Spark 作业的两种方式
Spark 作业不能直接嵌在 SpringBoot 进程里跑,因为两者对内存和 JVM 参数的要求完全不同。常见做法是两种:SpringBoot 通过ProcessBuilder调用spark-submit脚本,或者用 Quartz 定时扫描 HDFS 新数据目录后触发提交。
@Scheduled(cron = "0 30 1 * * ?") // 每天凌晨 1:30 执行 public void runDailyAnalysis() { String sparkSubmit = "/opt/spark/bin/spark-submit"; String jar = "/data/app/bike-analysis.jar"; String mainClass = "com.bike.analysis.DailyAnalysisJob"; String dt = LocalDate.now().minusDays(1).toString(); ProcessBuilder pb = new ProcessBuilder( sparkSubmit, "--class", mainClass, "--master", "yarn", "--executor-memory", "4g", "--num-executors", "8", jar, dt ); pb.redirectErrorStream(true); Process process = pb.start(); // 记录日志到文件 }用ProcessBuilder而不是 Java 代码里直接new SparkSession,原因在于 Spark 作业需要独立的 Driver 内存和 Executor 资源,和 Tomcat 容器共享 JVM 会导致 Full GC 互相影响。定时任务在凌晨跑日结,用户无感知,失败时通过日志文件排查。Cron 表达式0 30 1 * * ?的意思是每天 1:30,这个时间点通常是共享单车的低谷期,避免影响次日的统计分析结果。
5. 集群部署与内存配置:从单机 Demo 到分布式运行
很多人在本地跑通 Spark 后,一部署到服务器就各种报错。这一章梳理部署过程中最关键的参数和最容易踩的坑。
5.1 集群模式选择:为什么不用 local 模式
本地开发时用local[*]能直接跑通,但生产环境必须切换为 YARN 模式。YARN 模式的最大好处是资源统一管理和动态分配。在spark-defaults.conf里,有四个参数决定了作业的运行表现。
| 参数 | 推荐值 | 说明 |
|---|---|---|
| spark.executor.memory | 4g-8g | 单个 Executor 的内存,过大易触发 YARN 单容器限制 |
| spark.executor.cores | 2-4 | 单个 Executor 的 CPU 核数,避免过多线程竞争 |
| spark.driver.memory | 2g-4g | Driver 端内存,聚合结果大时调高 |
| spark.sql.shuffle.partitions | 100-400 | 影响 Shuffle 并行度,需随数据量调整 |
spark.executor.memory不是越大越好。YARN 默认的单容器最大内存是 8G,超过会被 ResourceManager 杀掉。我见过不少人把 Executor 内存调到 12G,结果作业一提交就报Container is running beyond physical memory limits。正确做法是先查集群的yarn.scheduler.maximum-allocation-mb,再按这个值的一半左右来设置 Executor 内存,留出 overhead 的空间。
5.2 串行提交与资源竞争问题
SpringBoot 定时任务如果同时触发多个 Spark 作业,而 YARN 队列资源有限,后面的作业会一直等待。推荐的思路是给不同作业设置不同的优先级,或者在 SparkConf 里强制开启 Fair Scheduler。
// 作业内部设置调度池,让核心作业优先生效 spark.sparkContext.setLocalProperty("spark.scheduler.pool", "production")调度池在 YARN 的capacity-scheduler.xml里配置。把日结报表放进production池,把实验性分析放进dev池,权重设为 2:1,这样即使同时提交,核心作业也不会被挤到后面。
5.3 数据倾斜的处理技巧
共享单车数据里,热点区域和头部用户的倾斜非常明显。处理倾斜有两个常用手段:两阶段聚合加随机前缀。
// 第一阶段:加随机前缀打散 val saltedDF = orderDF .withColumn("salt", (rand() * 10).cast("int")) .withColumn("salted_user", concat(col("salt"), lit("_"), col("user_id"))) // 用加盐后的用户 ID 做一次聚合 // 第二阶段:去掉前缀再做二次聚合(rand() * 10)生成 0 到 9 的随机整数,把每个大用户拆成 10 个小 key。这样原本集中在一个 Executor 上的数据被均匀分到 10 个任务。缺点是中间结果会多一层 shuffle,适合单个 key 数据量超过整体 10% 的场景。
6. 项目验证清单与运行检查:5 个必做的正确性检验
系统跑起来之后,怎么确认数据算对了,比确认不报错更重要。Spark 任务的常见问题是:流程跑完,结果集少了或多了,但日志提示 Success。这一章给出毕业设计答辩前必做的验证方法。
6.1 与原始数据的对账校验
每次分析任务结束后,写一个校验脚本,对比明细层的总数和聚合层的总数。
# 对比 HDFS 源数据的记录数和聚合结果的记录数 hdfs dfs -du -s /bike/order/2025/06/01/ spark-sql --master yarn \ --executor-memory 2g \ -e "SELECT COUNT(*) FROM bike_order WHERE dt = '2025-06-01';"du -s返回的是目录体积,COUNT(*)是明细层记录数。如果两者数量级差太多,优先检查分区过滤条件是否生效,尤其注意 dt 的格式必须一致,2025-06-01和2025-6-1在 HDFS 路径上对应不同目录,混用会导致读不到数据。
6.2 时间字段的时区一致性检查
时间字段是共享单车数据里最容易出错的地方。车锁终端上报的时间一般是本地时间,HDFS 存储时如果直接存字符串而不带时区,后续所有时间计算都会混乱。
# 用 Python 脚本抽查原始数据的时间字段格式 from pyspark.sql import SparkSession spark = SparkSession.builder.master("yarn").getOrCreate() df = spark.read.parquet("/bike/gps/2025/06/01/") df.select("report_time").distinct().show(10, truncate=False)检查report_time是否统一为yyyy-MM-dd HH:mm:ss格式,如果不统一,需要在写入明细层时用date_format标准化。日期函数在格式化时如果遇到NULL或空字符串,不会报错但会返回NULL,聚合结果会缺数据。
6.3 前端接口响应时间与字段映射验证
模拟前端调用接口,检查返回字段是否与 JSON 序列化后的类匹配。前后端联调时最容易出问题的是 Long 类型在 JavaScript 中精度丢失。
// 统一在返回前转成 String,避免 JS 精度丢失 public class BikeRideVO { private String orderId; private String distance; private String duration; }共享单车订单 ID 通常是雪花算法生成的 Long 类型,超过 JavaScript 的Number.MAX_SAFE_INTEGER后会丢精度,前端拿到的 ID 跟数据库对不上。所有 ID 类字段在 VO 层转成 String 是一种防御式做法。
6.4 任务失败重跑时的幂等性验证
定时任务失败后的重跑机制经常被忽略。Spark 写聚合表时如果用的是saveAsTable,重跑时表已经存在会报错,需要先删掉对应分区的数据再写入。
// 重跑时先把目标分区数据删掉,保证幂等 spark.sql("ALTER TABLE peak_hour_stats DROP IF EXISTS PARTITION (dt = '2025-06-01')")删除分区后重新计算,结果只包含当前批次的数据,不会重复累加。如果任务失败发生在写 HDFS 之后、写元数据之前,HDFS 上的孤儿文件会留在目标目录里,下次运行前需要清理。
6.5 用 curl 做接口层面的冒烟测试
部署完成后,用一条 curl 命令验证整个链路是否通畅。
curl -X GET "http://node01:8080/api/analytics/peak-hours?date=2025-06-01" \ -H "Authorization: Bearer $TOKEN" \ -w "\nHTTP_CODE:%{http_code}\n"-w参数输出 HTTP 状态码,状态码 200 只代表接口响应正常,不代表数据正确。还要对比返回里的order_count是否与 Spark SQL 手动查询的结果一致。如果返回空数组,检查 MySQL 的peak_hour_stats表数据是否被正确的 Spark 任务写入,以及表名是否和实体类映射一致。TOKEN从登录接口获取,放在环境变量里避免泄露。
本文还有配套的精品资源,点击获取