ARTICLE DETAIL

资讯详情

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

SpringBoot集成Spark:共享单车数据存储系统设计与避坑实践

SpringBoot集成Spark:共享单车数据存储系统设计与避坑实践 简介一份面向计算机专业毕业设计的共享单车数据存储系统项目基于 SpringBoot 与 Spark 技术栈覆盖论文、后端、前端与数据库脚本适合 Java 后端、大数据处理方向的课程设计与毕设参考。系统围绕单车骑行数据的采集、存储、分析与查询展开并结合分布式存储思路处理规模化数据。压缩包内共 374 个文件体积约 17.8MB主要包含 70 个 Java 后台类、35 个 Vue 前端组件、9 个 XML 配置、论文 docx 文档、SQL 脚本及安装运行批处理文件目录结构清晰便于按模块查阅。已有 51 人学习下载通过本套资料可理解 SpringBoot 构建 RESTful 接口、Spark 对接 HDFS 的数据处理方式并掌握前后端分离项目从编译到部署的完整流程。配套论文还能帮助梳理论文写作与系统设计思路适合需要快速构建同类型系统或学习大数据存储实现的读者。1. 共享单车数据存储SpringBoot 与 Spark 分工后的系统边界在哪一个共享单车平台的骑行订单高峰期一小时就是几十万条一天下来累计几千万条明细。把这些数据全部塞进 MySQL单表过亿之后聚合查询就开始拖垮业务库但如果只靠 SpringBoot 做存储层又解决不了批量清洗、按时间分区的数仓建模问题。标题里这套“springboot基于 Spark 的共享单车数据存储系统”本质上是把两条链路分开Spark 负责把原始骑行记录洗成可分析的分区表SpringBoot 负责把分析结果以接口形式暴露出去。这个方向在 Java 毕设和中小型团队的数仓改造里都很常见适合手里有批量历史数据、又需要对外提供查询接口的场景。这套系统的价值不在“用了大数据框架”而在它划清了一条边界明细数据归 Spark 管查询接口归 SpringBoot 管两者不抢数据源。下面从数据模型、入库脚本、接口封装和故障排查四段把这个方案完整走一遍。2. 数据落在哪共享单车原始数据的 HDFS 目录设计与 Parquet 分区表2.1 共享单车原始数据长什么样先看字段再定存储先把手里的骑行订单数据摊开看。常见导出格式是 CSV字段一般包含订单号、车辆编号、用户 ID、开锁站点编号、关锁站点编号、开锁时间、关锁时间、骑行时长秒、骑行里程米、支付金额元。有些平台还会带经纬度坐标但那通常是轨迹明细和订单明细要分开入库。这批数据有两个特征直接影响存储选型第一只追加、不修改一辆车的一次骑行结束之后就再也不会更新第二查询模式以时间聚合为主比如“某站点一天有多少辆车被骑走”“早晚高峰哪个区域骑行量最大”。这两个特征决定了它不适合直接丢给 MySQL 做主存储更适合放到 HDFS 上用 Parquet 列式存储按天做分区。字段层面要提前处理两件事时间字段统一成yyyy-MM-dd HH:mm:ss字符串或timestamp类型站点编号、车辆编号这类维度字段保持字符串避免后续 join 时频繁做类型转换。这个决定会影响所有下游查询越早统一越省事。2.2 Parquet 分区表为什么比 JSON 和 MySQL 更适合这套系统Parquet 是列式存储同一列的数据在物理上连续存放。做“按站点 group by 求和”这类聚合时Parquet 只需要读取站点、时间、金额三列不需要像 JSON 那样把整行记录全部读进内存再丢弃无用字段。分区设计推荐按天分区表结构长这样CREATE TABLE IF NOT EXISTS ods_bike_order ( order_id STRING, bike_id STRING, user_id STRING, start_station STRING, end_station STRING, start_time TIMESTAMP, end_time TIMESTAMP, duration_sec BIGINT, distance_m DOUBLE, amount DOUBLE ) USING PARQUET PARTITIONED BY (dt STRING);这里dt是分区字段存“日期字符串”推荐格式2024-06-01。这样用 Hive 风格的表结构建表SparkSQL 可以直接读写。分区字段单独建一个不要塞进大字段里否则每次查询都全表扫描。建表之后每天的数据写进当天分区查询时带WHERE dt 2024-06-01就能只扫一个分区。实测常见场景下带分区裁剪的查询和不带分区裁剪的查询性能差距能到几十倍。2.3 目录分层ODS、明细层和聚合层分别放什么共享单车数据不建议只建一张表。常见做法是三层目录ODS 层存放原始 CSV 转换后的 Parquet 数据理论上和源文件字段一一对应DWD 层做清洗过滤掉骑行时长为负、站点编号为空的脏数据ADS 层存放按站点、按小时聚合好的结果表供 SpringBoot 接口直接查询。/data/bike/ods/ # 原始订单明细按天分区 /data/bike/dwd/ # 清洗后明细订单号去重 /data/bike/ads/ # 站点日维度、小时维度聚合结果这个层级不是必须的但建议至少拆 ODS 和 ADS 两层。做过一次就会体会到如果只留一张原始表SpringBoot 做接口查询时每个请求都要全量扫数据拆出 ADS 层之后接口查询只需要碰几百行聚合结果压力小一个量级。这个目录结构同时承担了 HDFS 上的物理路径和 SparkSQL 表名映射两层含义表建在哪个库、路径挂在哪个目录下要和spark.sql.warehouse.dir保持一致。3. Spark 批量入库从 CSV 到分区表的入库脚本与四个调参项3.1 用 PySpark 写清洗逻辑过滤脏数据、去重、统一时间格式数据落地的第一步是把 CSV 源文件读进来做基本清洗。下面是一段常见的 PySpark 批处理脚本跑在 Spark 集群上每天把当天 CSV 写入 ODS 表from pyspark.sql import SparkSession, functions as F spark SparkSession.builder \ .appName(bike_order_etl) \ .config(spark.sql.shuffle.partitions, 200) \ .enableHiveSupport() \ .getOrCreate() df spark.read.csv( hdfs:///data/input/bike/2024-06-01/*.csv, headerTrue, inferSchemaTrue ) df df.filter( (F.col(duration_sec) 0) (F.col(start_station).isNotNull()) (F.col(end_station).isNotNull()) ) df df.withColumn(dt, F.to_date(F.col(start_time), yyyy-MM-dd HH:mm:ss)) df.write \ .mode(overwrite) \ .partitionBy(dt) \ .format(parquet) \ .saveAsTable(ods_bike_order)spark.sql.shuffle.partitions控制的是 shuffle 阶段的分区数默认值 200 在集群资源偏小时容易造成大量小文件改成和 executor 核数相关、或者按数据量估算的值入库速度会明显变化。partitionBy(dt)决定 HDFS 上按天分目录后续查询和增量写入都以这个字段为准。注意这里mode(overwrite)只覆盖动态分区涉及到的分区目录不会把整张表清掉这个行为在 Spark 3.x 里由spark.sql.sources.partitionOverwriteMode控制建议设成DYNAMIC。3.2 Spark 集群搭建后必调的四个参数网上讲 spark 集群搭建的教程很多但环境搭起来之后跑批翻车的十有八九是内存参数没配对。以下四个参数我每次都会显式写在spark-submit里不依赖集群默认值spark-submit \ --master yarn \ --deploy-mode client \ --driver-memory 4g \ --executor-memory 8g \ --executor-cores 4 \ --num-executors 10 \ --conf spark.sql.shuffle.partitions200 \ --conf spark.sql.sources.partitionOverwriteModeDYNAMIC \ bike_etl.py--executor-memory不是越大越好要给 YARN 的 overhead 留出余量同时小于节点物理内存的一半才不会把 NodeManager 拖垮。--executor-cores 4配合 8g 内存比较均衡如果每个 executor 给 8 个核但内存只有 8g执行时 GC 会明显变长。spark.sql.shuffle.partitions200对单日几百万到上千万条记录够用如果源数据超过 5 亿条建议按总数据量 / 每分区500万条的估算往上调。这些参数属于“玄学”成分最多的部分但也是 spark 内存紧张时最先要查的地方。3.3 增量追加与调度别每次跑批都全量重刷共享单车数据是持续产生的每天都会有新的 CSV 文件进入 HDFS 的 input 目录。如果每次跑批都读全量历史再覆盖写全表随着时间推移作业会越来越慢而且每次都产生全量新文件小文件数量会涨到 NameNode 告警。正确的做法是每天只处理当天新到的文件按dt分区写入写入模式用overwrite或append取决于当天文件是否可能重复生成。实际项目里我一般让调度平台Airflow 或 shell cron 都行在每天凌晨执行一次入库脚本脚本第一步先打印读取到的文件列表和行数确认源数据就位再开始清洗。这样第二天发现源文件有问题时还有后悔药可吃——直接删掉当天分区重跑即可不影响其他分区。增量跑批之后ADS 聚合层也按天刷新INSERT OVERWRITE TABLE ads_bike_station_daily PARTITION(dt) SELECT ... FROM ods_bike_order WHERE dt 当天 GROUP BY start_station, dt。这个查询把 ODS 层的几百万行压缩成几百行SpringBoot 的接口压力从此只跟 ADS 层大小挂钩。4. SpringBoot 查询服务把 Spark SQL 查询方式包成 REST 接口的三种方案4.1 三种对接方式嵌入 SparkSession、走 Thrift Server、提交批任务SpringBoot 和 Spark 通信不是只有一种姿势从业界落地经验看主要分三条路。第一种在 SpringBoot 进程里直接启动一个 SparkSession 当嵌入式计算引擎。这种方式开发最快查询接口里直接写spark.sql(SELECT ...)拿 DataFrame。但要注意这条路径的边界SparkSession 初始化时会预占数 GB 内存且不适合高并发请求适合数据量可控、以演示和毕设答辩为目标的项目。第二种启动独立的 Spark Thrift ServerSpringBoot 通过 JDBC 连接它执行 SQL。这种方式把 Spark 当数仓用SpringBoot 只写 JDBC 代码结构干净、并发能力比第一种强很多。代价是环境里要额外部署 Thrift Server连不上时排错路径长一点。第三种SpringBoot 通过 REST API 向 YARN 提交 Spark 批任务任务结果落到 HDFS/ADS 表SpringBoot 再查 ADS 结果表返回给前端。这种方式适合离线分析场景接口实时性要求不高时用。对共享单车存储系统这个题目第二种最贴近“数据存储系统”的定位既保留了 SparkSQL 的分析能力又让 SpringBoot 保持轻量。4.2 用一个预热的 SparkSession 提供查询接口如果选择第一种方式不要在 Controller 里new SparkSession——SparkSession 初始化有秒级开销一次请求就初始化一次接口基本没法用。正确做法是项目启动时统一创建存在一个 Spring 管理的 Bean 里后面所有查询共用。Component public class SparkSessionHolder { private SparkSession spark; PostConstruct public void init() { spark SparkSession.builder() .appName(bike-query-service) .master(yarn) .config(spark.sql.shuffle.partitions, 50) .enableHiveSupport() .getOrCreate(); } public SparkSession getSpark() { return spark; } }RestController RequestMapping(/api/bike) public class BikeStationController { Resource private SparkSessionHolder holder; GetMapping(/station/daily) public ListMapString, Object dailyStationSummary( RequestParam String startDate, RequestParam String endDate) { String sql SELECT start_station, dt, COUNT(*) cnt FROM ods_bike_order WHERE dt startDate AND dt endDate GROUP BY start_station, dt ORDER BY cnt DESC LIMIT 100; return holder.getSpark().sql(sql) .collectAsList() .stream() .map(row - { MapString, Object map new HashMap(); map.put(station, row.getString(0)); map.put(dt, row.getString(1)); map.put(cnt, row.getLong(2)); return map; }).collect(Collectors.toList()); } }这段代码里查询直接打在 ODS 表上走 Spark SQL 的谓词下推Spark 会把dt的过滤条件下推到文件扫描层。查询结果数据量控制在 100 条以内避免 DataFrame 大结果集反序列化带来的性能问题。注意 Controller 里拼 SQL 时做过过滤就足够安全但更规范的做法是加一层参数校验日期格式正则不过直接返回 400。4.3 查询走 spark sql 时的类型坑分区字段比大小、维度字段没收敛SparkSQL 接口写起来快但两个坑要提前预防。第一个是分区字段类型不一致。如果dt字段在表里是 string你传参数时用了dt 2024-06-01没问题但如果你建表时把dt建成了 date 类型参数传2024-06-01也能跑可一旦你写成dt 20240601分区裁剪直接失效全表扫描。所以建表时统一dt为 string查询参数永远用yyyy-MM-dd格式。第二个坑是站点编号这类维度字段数据不收敛。共享单车站点有可能在运营中新增、合并、下线如果汇总 SQL 没有做站点维表 join结果里会出现大量已废弃站点编号对着零数据。常见做法是在 ADS 层聚合时 join 一张dim_bike_station维表把不在维表里的站点编号归到unknown避免接口返回脏维度。5. 共享单车数据存储系统避坑指南集群 OOM 与分区倾斜的 6 条故障记录5.1 现象本地 IDEA 跑通部署到 spark 集群就 OOM本地开发环境的数据量通常只有几万条IDEA 里跑SparkSession用local[*]模式一切正常。部署到 spark 集群搭建的 YARN 环境后同样代码读全量历史数据作业跑一会儿就报Container killed by YARN for exceeding memory limits。原因有两个一是源数据量级变大了但 executor 默认内存只有 1g根本装不下 shuffle 数据二是spark.sql.shuffle.partitions没配默认 200 个分区同时落盘每个分区的临时文件叠加起来远超容器内存。解决方法是先把--executor-memory给到节点物理内存的一半左右再调大spark.sql.shuffle.partitions把单个分区的数据摊薄。这套配置不要写在代码里写在spark-submit命令行上因为不同环境节点规格不一样写死反而坑到下一个人。5.2 现象跑批作业越跑越慢日志里全是小文件告警系统上线初期每天入库增量数据一切正常。跑了一个月之后相同的数据量入库时间从 10 分钟涨到 40 分钟HDFS 上/data/bike/ods目录下出现成百上千个几十 KB 的小 Parquet 文件。原因spark.sql.shuffle.partitions设得太大而且每次INSERT按分区落盘时200 个 reducer 每个都往同一个分区目录写一个小文件。文件越碎下一次全表查询时 NameNode 的元数据开销越大。解决把partitionBy(dt)改成双层分区partitionBy(dt, hour)让每个分区的数据量更大同时把spark.sql.shuffle.partitions从 200 降到和 executor 总数匹配比如 10 个 executor 各 4 个核就设 40 或 80。如果已经产生大量小文件用ALTER TABLE ... CONCATENATE合并 Parquet 文件或者直接对 ODS 分区重写一次。5.3 现象按站点聚合热门站点的分区数据量是其他站点的百倍共享单车的潮汐效应会在数据层面直接暴露一个地铁口站点一天几千次开锁郊区站点一天三次。Spark 按start_station做 group by 时同一个站点所有数据都进同一个 reducer该 reducer 的数据量是其他 reducer 的百倍以上表现为整个作业卡在最后一个 stage。原因就是数据倾斜本质上是键分布不均匀。解决这个问题的常规套路是给热点 key 加盐先按start_station分组统计给超过阈值的站点编号拼接随机后缀把一个大键拆成多个小键聚合完成后再去掉后缀二次聚合。对共享单车这个场景还有一个更省事的方案——改用双层分区dt hour先做时间维度的预聚合把数据粒度变粗之后再做站点聚合倾斜程度会大幅缓解。5.4 现象SpringBoot 接口偶发超时日志里出现 SparkSession 不可用用嵌入式的 SparkSession 跑查询接口高峰期并发上来之后偶发org.apache.spark.SparkException: Job aborted重启 SpringBoot 服务又好了但过一阵又出现。原因是 SparkSession 不是线程安全的多个请求同时提交 SQL 时底层 SparkContext 的调度状态会互相干扰。不要在一个 Controller 里并发调用同一个 SparkSession这不是代码 bug是设计边界问题。解决在 SpringBoot 里对查询接口做串行化用synchronized或单线程线程池包住 Spark 查询调用或者干脆走 4.1 方案二把查询任务迁到 Spark Thrift Server 上SpringBoot 只维护 JDBC 连接池。对于毕设级别的系统用单线程线程池包一层就行实测一天十万次以内的查询压力完全扛得住。5.5 现象入库去重后数据反而变少了和源系统对不上账系统上线后做数据对账发现 ODS 层的记录数比源系统导出的总记录数少了几万条。逐日排查发现每天跑批时用了dropDuplicates([order_id])但部分订单在前一天已经被写入了——CSV 源文件会跨天重复推送同一个订单。原因源系统不是严格按订单完结时间导出的而是按“数据落库时间”导出导致一个订单可能出现在两个批次的 CSV 文件里。对账时如果只按单日去重就会把前一天已入库的订单洗掉。解决去重逻辑里增加时间窗口条件只对当天分区内做order_id去重不要跨分区去重或者建表时用order_id做桶字段按order_id分桶加sortBy(start_time)入库时按桶内去重。对共享单车这种以时间为维度的数据跨天重复是常态去重一定要限定分区范围。5.6 现象Spark SQL 查询结果和报表系统对不上差在时区ADS 层按小时聚合作业在每天凌晨跑聚合出的“0 点时段”数据比报表系统里的少一些。最后查到是时区问题源 CSV 里start_time是yyyy-MM-dd HH:mm:ss字符串但 Spark 读取时默认按服务器时区解析集群节点统一 UTC 时区时to_date转出的日期比实际早 8 小时。解决在 SparkSession 初始化时显式设置spark.conf.set(spark.sql.session.timeZone, Asia/Shanghai)或者入库时直接按字符串substring(start_time, 1, 10)截取日期作为分区字段绕开时区解析。这个坑只在集群时区和业务时区不一致时出现但一旦出现就是整表数据错位恢复成本很高所以建议建表 schema 里时间字段全用字符串业务侧再统一转换。6. 分区裁剪与执行计划验证把 5 亿条骑行记录的聚合查询从分钟级拉到秒级最后分享一个我每次跑完入库都会做的验证技巧用EXPLAIN看执行计划确认查询真正走的是分区裁剪而不是全表扫描。方法很简单Spark SQL 里直接执行EXPLAIN SELECT start_station, COUNT(*) FROM ods_bike_order WHERE dt 2024-06-01 GROUP BY start_station;如果执行计划里出现Partition Filters: [dt2024-06-01]说明分区裁剪生效如果看到FileScan parquet后面没有分区过滤条件那就要检查dt字段类型是不是和查询参数不一致。这个检查 30 秒就能做完能挡住 80% 的“为什么查询这么慢”类问题。另外一个习惯是容量验证。我会在系统上线前用脱敏数据模拟一个月的骑行记录约 5 亿条先量出单日分区的 Parquet 文件总大小再估算一年的 HDFS 占用。共享单车的订单数据量单日分区做得好大概在 200MB 到 500MB 之间一年下来也就 100GB 到 200GB 的量级三节点的小集群完全顶得住。这个估算能给团队一个明确预期这套系统不需要上十几台机器的大集群硬件成本可控做起来不虚。关于接口层我还有一个个人教训SpringBoot 查询接口返回前一定要对 DataFrame 结果做LIMIT限制。曾经有一次没加LIMITcollectAsList()直接把几百万行聚合结果拉进 JVM 堆内存接口直接 OOM服务重启才恢复。从那以后凡是 Spark 查询结果转 List我有两条铁律SQL 里必须带LIMITJava 侧再限制collectAsList()的规模。这套 SpringBoot 加 Spark 的方案适合数据量在千万到亿级、查询以离线聚合为主的场景。它比纯 MySQL 方案跑得快比纯 Spark 方案好对接业务系统边界就在“离线批处理 在线查询”这条线上。把 ODS、ADS 分层做好把分区字段类型统一把内存参数按实际数据量调一遍这个系统维护成本很低。希望帮到你。本文还有配套的精品资源点击获取
返回列表