ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

共享单车大数据存储:Spark批处理与Spring Boot集成实践

共享单车大数据存储:Spark批处理与Spring Boot集成实践 简介一套基于SpringBoot与Spark的共享单车数据存储系统毕业设计论文与源码包面向计算机科学与技术、人工智能等专业学生也适合作为课程作业或大数据实践项目系统围绕共享单车数据的采集、分布式存储、Spark分析与前端交互展开能够帮助学习者掌握分布式存储、大数据计算及前后端开发等技术。资源共374个文件以Java后端、Vue前端、XML/SQL脚本及Word论文为主另含JS、CSS、图片、配置文件等静态与支撑资源压缩包约17.8MB内部目录结构清晰并提供安装、运行、构建脚本便于快速导入项目并验证功能。目前已有50人学习浏览。除完整项目源码和论文外包内还包含数据库初始化脚本、前后端配置和批处理启动脚本读者可结合论文理解共享单车数据管理的业务场景与系统架构参考其中的数据表设计、Spark分析思路和RESTful API实现在此基础上扩展调度优化、用户画像等新功能。1. 共享单车数据存储系统为什么要把 Spark 和 Spring Boot 搭在一起共享单车系统最麻烦的不是还车逻辑而是每天几千万条的轨迹和订单流水。这些数据单独看都没用但按车辆、站点、时间片聚合以后就能算出城市骑行热力图、潮汐调度和车辆损坏预测。把清洗、聚合、存储交给 Spark把对外接口和任务提交交给 Spring Boot是后台开发里比较常见的组合。下面会从数据模型开始讲到 Spark 作业怎么写、Spring Boot 怎么触发任务最后给参数调优和验证方法。适合正在做毕业设计选型的开发者也适合想在生产环境里把这两个框架拆开部署的工程师。2. 系统分层与数据流Spring Boot 负责调度Spark 负责算共享单车数据的处理链路可以拆成三个步骤收集原始数据、批量清洗聚合、对外提供查询。这三个步骤对应的技术栈各不相同Spring Boot 和 Spark 的组合本质上就是把“服务态”和“计算态”分开避免两个框架互相挤占资源。2.1 先分清批处理和实时处理Spark Streaming 要现在引入吗共享单车数据有两种常见的消费方式。一类是实时告警比如车辆驶出电子围栏、订单连续超时、设备心跳丢失这类场景要求秒级响应通常用 Spark Streaming 或 Flink 配合消息队列完成。另一类是批量统计比如每辆车的日均骑行时长、每个站点的借还峰值、周末和工作日骑行量对比这类场景需要扫描全量数据用 Spark 核心的 DataFrame 操作就够了。实际做这个系统时我倾向于先把批量统计做透再考虑实时链路。原因很简单共享单车运营中真正需要毫秒级计算的规则并不多绝大多数“大屏实时看板”用 Spark Streaming 消费 Kafka 算出来的结果和定时任务每五分钟刷一次 Redis 相比业务上的差距并不大。而引入 Streaming 会带来消息队列、消费位点、窗口延迟这一堆额外复杂度在毕设或中小型项目中很容易出现 Kafka 集群还没搭好截止日期先到了。这背后的选型原则可以记成一句话能离线算的不要实时算能聚合查的不要扫明细。实时化改造应该发生在“批处理结果已经不够用”之后而不是系统设计第一天就铺满。共享单车这个选题先做好 Spark 批处理已经能覆盖 80% 的数据存储需求。2.2 存储层怎么设计原始数据进 HDFS/MinIO聚合结果进 MySQL原始数据包括单车 GPS 上报、订单流水、用户扫码记录。这些数据的读取模式是“写一次读多次”基本不会单条修改所以适合用文件存储。Spark 对 Parquet 和 ORC 的列式存储支持得比较好查询时只读需要的列性能比直接读 JSON 高很多。聚合结果的数据量会骤减。以每天两千万条轨迹为例按车辆聚合后可能只剩几万条统计记录完全放得进 MySQL。Spring Boot 的 Web 接口只查这一层就不需要和 Spark 抢资源。常见的硬件条件下可以用 MinIO 代替 HDFS因为它部署简单一个 Docker 命令就能起来同时兼容 S3 协议Spark 可以通过s3a://直接读取。下面是一条启动命令docker run -d -p 9000:9000 -p 9001:9001 \ -v /opt/minio/data:/data \ minio/minio server /data参数说明9000是 S3 API 端口9001是管理控制台端口。在服务器上使用时/opt/minio/data要挂载到单独的磁盘避免和系统盘 IO 争抢。Spark 作业里访问 MinIO 需要设置spark.hadoop.fs.s3a.endpoint和spark.hadoop.fs.s3a.access.key对应的 secret key 写在作业提交的配置文件中。下面用一张表把存储分层列清楚层级数据内容存储介质读取方式原始层轨迹点、订单流水MinIO / HDFSSpark 批读明细层清洗后的骑行明细Parquet 文件Spark SQL聚合层车辆日统计、站点热度MySQL / PostgreSQLSpring Boot 查询缓存层高频接口结果RedisSpring Boot 读取这个分层的核心好处是职责分离Spark 永远不直接读业务库Spring Boot 也永远不处理原始大文件。两边的扩展可以独立进行Spark 集群计算能力不够时只加节点就行不影响 Web 服务的稳定性。2.3 Spring Boot 和 Spark 怎么连SparkLauncher 是更稳的姿势Spring Boot 和 Spark 的连接方式直接决定系统上线后的稳定性。第一种方式是在 Spring Boot 应用里直接 new 一个 SparkSession比如SparkSession.builder().master(local[2]).getOrCreate()。这种方式在开发环境没问题代码写起来最顺手但一旦部署到服务器Web 服务和 Spark Driver 会挤在同一个 JVM 里。接口并发一高Full GC 或 OOM 就会把 Spark 作业一起带走。另一种方式是用 SparkLauncher 提交独立的 Spark 作业Spring Boot 只负责拼参数、调用启动器、查询提交状态。作业在 YARN 或者 Standalone 集群上运行Driver 和 Executor 都在独立进程里Spring Boot 重启也不影响正在跑的任务。共享单车系统里的统计任务一般要跑几分钟这种物理上的进程隔离非常值得。如果只是在个人电脑上用 IDEA 做演示Spring Boot 里直接建一个local[2]的 SparkSession 也能接受但要清楚它的上限在哪里。这里单独提一下 Spark 集群搭建。如果你用的是 Standalone 模式部署时记得把每个 Worker 的内存和 CPU 核数写清楚并确认spark-env.sh里没有乱配SPARK_WORKER_CORES否则可能出现集群资源看起来很大、但单个作业只能拿到 1 个核的错觉。用 YARN 的话则要提前定好队列避免和其他任务抢资源。第 4 章会给提交代码这里先记住结论调度和服务要轻计算和存储要独立。3. 共享单车数据模型与 Spark 作业核心实现数据模型设计是 Spark 作业能不能写好、Spring Boot 查询快不快的前提。共享单车系统的模型并不复杂难在字段多、量大、聚合口径经常变。所以模型设计要遵循“原始层宽一点聚合层瘦一点”的原则把需要做判断的字段尽量在原始层保留把接口要用的统计值收敛到聚合层。3.1 订单表和轨迹表的字段怎么定义才够用订单表里尽量保留原始字段因为后面做特征分析时还不知道哪些字段会被用到。一般包含订单号、车辆编号、用户编号、借车站点 ID、还车站点 ID、借车时间、还车时间、骑行费用、骑行时长、骑行里程。轨迹表则更宽每一条 GPS 上报都算一条记录包含轨迹点 ID、订单号、经度、纬度、速度、电量、设备时间、上报时间。轨迹表的设计重点不在字段而在分区。按date字段做分区写代码时指定date 2025-06-01Spark 的 Predicate Pushdown 会直接跳过无关分区查询速度提升非常明显。原始轨迹文件在落盘时尽量按照date目录组织比如/raw/order/2025/06/01/这样 Spark 读路径时天然按目录裁剪。这里给出 MySQL 中聚合表的建表语句该表存的是每天每辆车的结果CREATE TABLE order_daily_stat ( id BIGINT PRIMARY KEY AUTO_INCREMENT, biz_date DATE NOT NULL, bike_id VARCHAR(32) NOT NULL, order_cnt INT NOT NULL DEFAULT 0, total_duration_min DECIMAL(10, 2) NOT NULL DEFAULT 0, total_distance_km DECIMAL(10, 2) NOT NULL DEFAULT 0, UNIQUE KEY uk_biz_bike (biz_date, bike_id) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;参数说明uk_biz_bike联合唯一索引是幂等写入的关键Spark 重复执行同一天的作业时可以靠它避免出现双份数据。duration和distance用DECIMAL而不是FLOAT避免在累加时产生精度漂移。注意bike_id不要用数字类型不同运营商的共享单车编号往往带字母前缀用VARCHAR(32)最稳妥。3.2 用 Spark 清洗和计算订单数据的完整 Java 代码Spark 作业的入口类建议只做四件事读数据、清洗、聚合、写出。下面这段代码是订单日统计作业的核心部分使用 Spark 3.x 的 DataFrame API 编写。import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; import static org.apache.spark.sql.functions.*; public class BikeOrderDailyJob { public static void main(String[] args) { if (args.length 3) { System.err.println(Usage: BikeOrderDailyJob inputPath jdbcUrl tableName); System.exit(1); } String inputPath args[0]; String jdbcUrl args[1]; String tableName args[2]; SparkSession spark SparkSession.builder() .appName(BikeOrderDailyJob) .config(spark.serializer, org.apache.spark.serializer.KryoSerializer) .getOrCreate(); DatasetRow df spark.read() .option(header, true) .option(inferSchema, true) .csv(inputPath); DatasetRow cleaned df .filter(col(start_time).isNotNull() .and(col(end_time).isNotNull()) .and(col(bike_id).notEqual(null))) .withColumn(start_date, to_date(col(start_time))) .withColumn(duration_min, round(col(duration_second).cast(double).divide(60), 2)); DatasetRow dailyStat cleaned .groupBy(bike_id, start_date) .agg(count(order_id).alias(order_cnt), sum(duration_min).alias(total_duration_min)); dailyStat.write() .mode(append) .option(rewriteBatchedStatements, true) .jdbc(jdbcUrl, tableName, buildProperties()); spark.stop(); } }逻辑说明read().csv()返回的是 DataFramefilter里的几个条件把开始时间为空、结束时间为空、车辆编号为null的脏数据排掉。withColumn新增了两个字段一个是日期一个是把秒转换成分钟的时长。groupBy按车辆和日期聚合用count和sum得到订单数和总时长最后用append模式写入 JDBC。参数说明inferSchematrue会自动识别类型但会在读文件时额外扫描一遍文件很大时可以改成预定义schema来提速。rewriteBatchedStatementstrue是 MySQL JDBC 驱动的重要参数能合并多条 insert 语句批量插入速度能提高一个数量级。jdbcUrl可以在代码里通过useSSLfalseserverTimezoneAsia/Shanghai固定下来避免和服务器时区不一致导致时间错乱。3.3 数据倾斜场景下的分组优化共享单车订单在时间和空间上都不均匀。早高峰地铁站周边的订单量是平峰时间段的几十倍直接groupBy(bike_id)会让少数热点车辆所在的任务特别慢。Spark 的调度是短板效应整个作业等最慢的一个任务所以要把热点 key 打散。一种常见的优化是加盐给bike_id拼一个随机前缀先做局部聚合再按原始 key 聚合。比如withColumn(salt, floor(rand(seed) * 100))然后groupBy(salt, bike_id, start_date)算完后去掉salt再聚合。这样热 key 的负载被分散到多个 Executor 上代价是会多一次 Shuffle。数据量在亿级以下时直接调大spark.sql.shuffle.partitions配合 Adaptive Query Execution 往往更省事。如果不想改造作业逻辑可以在 Spark 配置里开启以下参数spark.sql.adaptive.enabledtrue spark.sql.adaptive.coalescePartitions.enabledtrue spark.sql.adaptive.skewJoin.enabledtrue这三个参数在 Spark 3.x 已默认开启但老集群可能被运维关掉提交作业时显式加上更保险。开启后Spark 会在运行时检测数据倾斜的 Join 和 Reduce把大分区拆小小分区合并减少不必要的任务开销。对于共享单车的大量轨迹文件这个机制能明显降低落盘时间。4. Spring Boot 集成 Spark 的接口实现与配置Spring Boot 这边要做的事情比 Spark 作业简单但坑不少。最常见的问题是依赖冲突和任务提交时机的选择。这一章的落地顺序是先配置依赖再写提交服务再写 REST 接口最后处理幂等。4.1 创建 Spring Boot 项目时如何避免和 Spark 打架用 IDEA 创建 Spring Boot 项目很简单但引入 Spark 依赖时常常会踩依赖冲突的坑。Spring Boot 默认引入的 Jackson 版本和 Spark 内部的 Jackson 并不同步启动时会出现jackson-module-scala的NoSuchMethodError。最常见的解法是Spring Boot 项目里不依赖完整的 Spark Core只依赖spark-launcher。这个模块只负责提交任务不会把 Spark 的一整套类都引进来。如果必须在 Spring Boot 里直接处理 Spark DataFrame就用 Maven 的exclusions把jackson-module-scala排除掉或者把 Spring Boot 的版本调整到 2.6.x/2.7.x 这个区间配合 Spark 3.x 能少踩一些坑。太新的 Spring Boot 3.x 用了 Jakarta EESpark 3.3 之前的部分代码还基于 Java EE硬组合会增加额外适配工作。Spring Boot 自动装配也会扫描数据源配置如果作业 JAR 里也带了数据源容易产生重复初始化。所以我一般把 Spark 作业打成独立 JARSpring Boot 只保留 Web 侧的连接池。下面是一段可用的 Maven 依赖配置示例dependency groupIdorg.apache.spark/groupId artifactIdspark-launcher_2.12/artifactId version3.3.2/version /dependency注意这里写spark-launcher_2.12Scala 版本要和 Spark 集群一致。如果集群用的 Scala 2.12作业 JAR 也用 2.12依赖里就不要再混入 2.13 的包。版本号建议以实际部署环境为准代码示例只作为参考。4.2 用 SparkLauncher 提交独立任务的封装实现Spring Boot 侧的任务提交核心不是自己写线程池而是调用SparkLauncher。下面是一个完整的提交服务包含任务日志收集和超时处理。Service public class SparkSubmitService { private static final Logger log LoggerFactory.getLogger(SparkSubmitService.class); Value(${spark.home}) private String sparkHome; public String submitJob(String appJar, String mainClass, String date) { try { Path logDir Paths.get(logs).toAbsolutePath(); Files.createDirectories(logDir); String logFile logDir.resolve(spark_ System.currentTimeMillis() .log).toString(); SparkLauncher launcher new SparkLauncher() .setSparkHome(sparkHome) .setAppResource(appJar) .setMainClass(mainClass) .setMaster(yarn) .setDeployMode(cluster) .setConf(spark.executor.memory, 2g) .setConf(spark.executor.cores, 2) .setConf(spark.sql.shuffle.partitions, 200) .addAppArgs(/raw/order/ date, jdbc:mysql://localhost:3306/bike_store, order_daily_stat); Process process launcher.launch(); log.info(launched spark process, log at {}, logFile); return logFile; } catch (Exception e) { throw new IllegalStateException(spark job submit failed, e); } } }逻辑说明setSparkHome指定 Spark 安装目录setAppResource指向已经打好的作业 JARsetMainClass指向作业入口。setDeployMode(cluster)会让 Driver 在 YARN 的 ApplicationMaster 中运行避免 Spring Boot 进程变成 Driver。addAppArgs把业务参数传给main方法。这里故意没有调用waitFor()因为统计任务可能跑 10 分钟同步等待会把 Web 请求线程占满。更合理的方式是提交后立刻返回一个logFile路径由另一个查询接口读取该文件的最后几行用来判断作业是否卡住。如果需要严格的任务状态管理可以结合 YARN REST API 查询application_*的状态。4.3 Spring Boot 的 REST API 只查聚合结果对外暴露的接口层要尽量薄。共享单车系统里需要对外提供的高频接口主要是查询某辆车某段时间的骑行统计、查询某个站点的借还量、触发某天数据的重算任务。RestController RequestMapping(/api/bike) public class BikeStatController { Autowired private OrderDailyStatRepository statRepository; Autowired private SparkSubmitService submitService; GetMapping(/daily) public ListOrderDailyStat getDaily(RequestParam String bikeId, RequestParam String startDate, RequestParam String endDate) { return statRepository.findByBikeIdAndBizDateBetween(bikeId, startDate, endDate); } PostMapping(/rerun) public String rerun(RequestParam String date) { if (!date.matches(\\d{4}-\\d{2}-\\d{2})) { throw new IllegalArgumentException(date format must be yyyy-MM-dd); } return submitService.submitJob(spark-jobs.jar, com.example.BikeOrderDailyJob, date); } }逻辑说明/daily接口直接查 MySQL 聚合表所有查询都走biz_date索引避免大范围扫描。/rerun接口接收日期参数先做正则校验防止任意字符串作为路径传入再通过submitService提交 Spark 作业返回日志文件路径。这样设计的好处是Spark 集群出问题时历史聚合结果仍然能正常查询只有重算任务会失败。系统把“读服务”和“计算任务”之间的互相影响降到了最低。4.4 写入幂等性处理先删当日再 appendSpark 作业重跑是常态。网络抖动、MySQL 锁表、YARN 容器被杀都会导致作业失败。如果在 Spark 写入时用overwrite会把整张表替换掉历史数据也可能被波及。更安全的方式是在 Spring Boot 提交任务前先通过 JdbcTemplate 删除当天的数据再让 Spark 以append模式写入。jdbcTemplate.update(DELETE FROM order_daily_stat WHERE biz_date ?, date); submitService.submitJob(spark-jobs.jar, com.example.BikeOrderDailyJob, date);说明这样 Spark 作业代码里完全不用关心删除逻辑Web 侧因为已经知道日期删除操作天然在作业启动前完成。如果作业失败当天数据缺一条都查不出来运营团队会第一时间发现问题。注意删除和提交之间没有原子性如果提交失败需要保留一个DELETE的补偿接口。5. 部署后的验证与调优先确认数据对再碰参数作业部署到集群以后第一反应不能是“跑通了就行”而是要确认输出是不是真的可信。共享单车系统里一个聚合值错了可能直接表现为某辆车的骑行时长比一天总秒数还大这类问题在 Spark UI 里看不出来必须靠对账查出来。5.1 先跑通最小案例再上集群本地开发时不要一上来就提交 YARN。常见做法是用local[*]模式加一份 1000 行的采样数据确认清洗和统计结果正确再打 JAR 提交到集群。这里的最小案例指的是把原始 CSV 截取前几个小时数据跑一次完整作业然后人工检查 MySQL 里的聚合结果是否合理。如果发现时间字段差 8 个小时多半是时区问题检查 JDBC URL 里的serverTimezone以及 Spark 读取 CSV 时的timestampFormat。这种问题在本地不明显因为 IDEA 所在机器和服务器时间往往不同。5.2 用对账 SQL 验证聚合结果可信拿一个已经跑完的日期分别统计原始 CSV 里的订单量和 MySQL 聚合表里的order_cnt总和。两者相等说明清洗阶段没有丢数据。SELECT SUM(order_cnt) FROM order_daily_stat WHERE biz_date 2025-06-01;如果这条记录小于原始文件行数很可能是轨迹文件里包含了大循环中的重复上报检查清洗阶段的去重条件是否放在groupBy之前。除了对总数还可以抽查几个热点车辆 ID对比单车的原始订单和聚合结果这样能把数据正确性的判断落到具体业务对象上。5.3 一个收益明显的配置开启 Kryo 序列化Spring Boot 和 Spark 都有自带的对象序列化默认的 Java 序列化在 Shuffle 时会多花很多内存。共享单车轨迹字段多一条记录可能包含时间、电量、经纬度等十几个字段Kryo 可以显著压缩序列化体积。在作业代码的 SparkSession.builder 中加上这几行.config(spark.serializer, org.apache.spark.serializer.KryoSerializer) .config(spark.kryo.registrationRequired, false) .config(spark.kryoserializer.buffer.max, 128m)参数说明registrationRequiredfalse让 Kryo 自动注册类避免ClassNotFoundbuffer.max建议设到 128m防止某个大对象序列化时缓冲区溢出。开启之后你会发现 Shuffle 的落盘量下降作业时间能减少约两成。对于共享单车这种字段结构固定的数据还可以实现自定义 Kryo Registrar 来进一步压缩但多数情况下默认配置已经够用。本文还有配套的精品资源点击获取
返回列表