免费获取学习方案
ARTICLE DETAIL

资讯详情

深耕编程基础知识与建站技术分享的一线实战洞察。

共享单车数据存储实战:Spark 清洗聚合 + Spring Boot 接口查询

共享单车数据存储实战:Spark 清洗聚合 + Spring Boot 接口查询 简介针对共享单车数据存储与处理的实际需求这份毕业设计资源整合了SpringBoot后端与Spark大数据处理框架完整实现了从数据采集、分布式存储到分析查询的全套系统并配套论文文档。包内共374个文件约17.8MB以70个Java源码、35个Vue前端页面、21个JavaScript脚本为核心辅以SQL数据库脚本、XML配置、vue组件及多份设计文档代码结构清晰便于二次开发与论文对照学习。已有51人学习下载项目经过严格测试验证可正常运行适合计算机科学与技术、人工智能等专业的学生用于毕业设计、课程项目或大数据方向入门实践。通过这套资料读者既能掌握SpringBoot搭建RESTful API与模块化开发的思路也能了解借助Spark对大规模单车数据进行统计预测、调度优化的具体方法是一份工程与学术结合度较高的参考资源。1. 共享单车数据存储为什么 Spark 和 Spring Boot 会出现在同一个标题里做共享单车业务的团队每天会积压上亿条轨迹点高峰期每分钟几万次写入MySQL 单库很快就撑不住。更麻烦的是运营要看的不是某一条订单而是「今天早高峰哪个区域的车辆密度最高」「哪个片区的车 3 小时没被骑过」这种查询一上聚合就慢得没法用。这个标题里的系统本质是把数据链路拆成两段Spark 负责把海量原始轨迹数据清洗、聚合、落成分析友好的存储格式Spring Boot 负责把处理好的数据用接口提供给上层应用。适合正在做毕业设计的学生也适合想给中小型项目补大数据处理能力的 Java 开发。你不需要一台真正的集群单机 Spark 就能跑通关键是理解数据在两端之间怎么流动。2. 先看数据和存储模型这东西到底在存什么为什么不能全塞 MySQL2.1 共享单车数据的三个典型特征共享单车数据不是单纯的一张订单表它至少包含三类数据订单记录、GPS 轨迹点、车辆状态变更。订单记录的数据量相对可控一天几十万条但轨迹点完全不同——车辆在骑行过程中每 5 到 10 秒上报一个点一台车骑 30 分钟就是 200 多个点一万辆车同时在跑就是几百万条记录一天下来轨迹点的规模在亿级。车辆状态数据又是另一种节奏开关锁、故障上报、电池电量变化事件频率高但单条体积小。这三类数据凑在一起给存储系统提出了三个要求写入要能扛住高并发时序数据查询要支持按时间范围和地理区域做聚合还要能容忍数据延迟——运营看的是「昨天全天」和「过去一小时」的统计不需要每一秒都精确。这些特征决定了存储不能只靠一张大表而是要把「原始明细」和「聚合结果」分开对待。这个标题里的共享单车数据存储系统解决的正是这种分层存储的问题。2.2 冷热分层HDFS 存原始轨迹MySQL 存聚合结果我一般把这个系统的存储分为两层热存储和冷存储。热存储负责给业务接口提供秒级响应存的是聚合后的结果数据放在 MySQL 里就够了。冷存储负责保存全部原始轨迹点放在 HDFS 上用 Parquet 格式按天和城市分区。为什么不把原始轨迹也放 MySQL简单算一笔账一天一亿条轨迹记录每条按 100 字节算一天就是 10 GB一个月 300 GBMySQL 不是存不下是查询时全表扫描的代价会让接口彻底失去响应。Spark 在这条链路里的角色是把原始数据从 Kafka 或者日志文件里读出来做清洗聚合再分别写到两层存储里。这样可以避免把「分析引擎」和「业务存储」混在一起。有人问过我要不要用 HBase 或者 ClickHouse 代替这个组合我的看法是如果你的团队对大数据生态不熟单机 Spark HDFS MySQL 是最容易跑通的组合后续要扩展也就是加节点的事不会推倒重来。ClickHouse 查询确实快但数据接入和更新机制对新手来说是个黑匣子排错成本高。2.3 表结构设计分区怎么设字段怎么留存储系统里的 MySQL 表要按查询方式来设计不是按业务对象来设计。运营端最常见的两个查询是「查某城市某天的区域车辆热力」和「查某辆车的轨迹回放」。所以 MySQL 侧我建议只保留两张表区域聚合表和订单摘要表。区域聚合表按城市 日期 小时 网格 ID 做主键预先算好车辆数、平均停留时长、骑行次数订单摘要表则保留每次骑行的起终点经纬度、里程和时长。轨迹明细不进 MySQL需要回放时直接从 HDFS 上读 Parquet 文件。HDFS 侧的 Parquet 文件用两层分区第一层 date第二层 city_id。分区的好处是查询时能直接剪掉无关目录避免 Spark 全表扫描。轨迹点表的核心字段包括 bike_id、order_id、lon、lat、speed、timestamp、bike_status。注意尽量把经纬度存成 DOUBLE不要用字符串不然后续地理计算的精度和速度都吃亏。CREATE TABLE bike_region_agg ( city_id INT, stat_date DATE, stat_hour INT, grid_id VARCHAR(32), bike_count INT, avg_ride_minutes DECIMAL(6,2), PRIMARY KEY (city_id, stat_date, stat_hour, grid_id) );这张表是 Spark 聚合后写入 MySQL 的目标表。grid_id 是地理网格编码把城市按 500 米 x 500 米划分可以用 Geohash 的前几位直接生成。主键设计成四个维度的组合是为了让同样的区域在同一个小时内只保留一条记录Spark 写入时用INSERT ... ON DUPLICATE KEY UPDATE做幂等这样即使批处理任务重跑也不会产生重复数据。2.4 为什么用 Spark 而不是纯 MapReduce 或 Flink选 Spark 是因为这个场景的批处理特征太明显。共享单车的运营分析是按天和按小时跑的没有复杂的窗口计算MapReduce 写起来啰嗦而且中间结果落盘的机制让调试变得很痛苦。Flink 的实时能力确实强但当数据从 Kafka 到 HDFS 这个链路只需要小时级延迟时引入流计算是在给自己增加运维负担。Spark 的 DataFrame API 能直接读 Kafka、写 Parquet一套代码既支持本地跑也支持上集群这是它在这个项目里最受欢迎的原因。3. 用 Spark 把原始轨迹清洗成可落库的数据核心作业怎么设计3.1 从零搭一个 Spark 作业最小能跑的框架无论是用 Scala 还是 Java 写 Spark 作业入口长得都差不多。这里给你一个 Scala 版本的骨架它从 Kafka 读共享单车轨迹原始数据做基本解析后落成 Parquet。建议用 Scala 是因为 DataFrame API 的单行表达式在 Scala 里写起来最干净Java 写同样的逻辑会多很多样板代码。object BikeShareStorageJob { def main(args: Array[String]): Unit { val spark SparkSession.builder() .appName(bike_share_storage) .config(spark.sql.shuffle.partitions, 200) .master(local[*]) .getOrCreate() import spark.implicits._ val rawDF spark.readStream .format(kafka) .option(kafka.bootstrap.servers, localhost:9092) .option(subscribe, bike-gps-topic) .option(startingOffsets, latest) .load() val parsedDF rawDF .selectExpr(CAST(value AS STRING) as json) .select( from_json($json, schema).alias(data) ) .select(data.*) val cleanedDF parsedDF .filter($lon.isNotNull $lat.isNotNull) .filter($lon 73 $lon 135) .filter($lat 18 $lat 53) cleanedDF .writeStream .format(parquet) .partitionBy(city_id, dt) .option(path, hdfs://localhost:9000/bike_data/trajectory) .option(checkpointLocation, /tmp/bike_ckpt) .trigger(Trigger.ProcessingTime(10 minutes)) .start() .awaitTermination() } }这段代码里最关键的两个参数是partitionBy和trigger。partitionBy(city_id, dt)决定了写入 HDFS 后的目录层级是城市下面挂日期这样后续按城市和天查询时Spark 只需要扫对应的目录能省掉大量 I/O。trigger设置成每 10 分钟一次批量写入既能够保证数据接近实时可见又不会因为频繁落盘把小文件撑爆。3.2 清洗规则这四类脏数据必须处理轨迹数据有多脏做了才知道。最常见的四类问题经纬度漂移——车辆停在原地但 GPS 跳出去几百米时间乱序——网络延迟导致先收到后发的点重复上报——同一秒收到两遍相同的记录还有空值字段。漂移点要是不清洗掉后面做区域聚合时会把车辆数算错运营看到某片区域忽然多出 30 台车实际一台都没有。清洗规则不要写得太复杂按优先级处理就好。第一优先是过滤空值和超范围经纬度用前面代码里filter条件就能做第二优先是去重用dropDuplicates(bike_id, order_id, timestamp)去掉完全重复的记录第三优先是漂移修正计算相邻两个轨迹点的速度如果速度超过 80 km/h 就丢弃前一个点。注意漂移检测不能只看速度还要看时间间隔间隔 10 秒的两点距离 500 米和间隔 1 秒的两点距离 500 米前者可信后者一定有问题。val withSpeedDF cleanedDF .withColumn(prev_lon, lag(lon, 1).over(Window.partitionBy(bike_id, order_id).orderBy(ts))) .withColumn(prev_lat, lag(lat, 1).over(Window.partitionBy(bike_id, order_id).orderBy(ts))) .withColumn(prev_ts, lag(ts, 1).over(Window.partitionBy(bike_id, order_id).orderBy(ts))) .withColumn(distance_m, haversine($lon, $lat, $prev_lon, $prev_lat)) .withColumn(time_diff_s, ($ts - $prev_ts) / 1000) val finalDF withSpeedDF .filter($time_diff_s 0 || $distance_m / $time_diff_s 80)这段漂移过滤的写法用的是窗口函数lag它能让每一条记录看到同一次骑行中上一条记录的位置。haversine是一个自定义 UDF 或内置函数用来计算两个经纬度点之间的距离。过滤条件是「速度不能超过 80 km/h」这是一个行业里常用阈值共享单车用户骑不了这么快超过一定是 GPS 漂移或者上报错误。如果你的业务场景里有电动滑板车可以把阈值放宽到 100但没有必要随意调大。3.3 聚合写入 MySQL用 Repartition 控制好并行度清洗完的明细数据落 HDFS 只是存储系统的一半还要产出聚合结果给 Spring Boot 查询。聚合逻辑是把轨迹点按时间窗口和地理位置分桶统计每个网格内的车辆数。这里有个常见的性能翻车点大量数据分布在少数几个热点网格比如地铁站周围导致 Shuffle 阶段数据倾斜少数几个 Executor 忙死其他 Executor 闲着。解决方案是给groupBy的 key 加一个前缀或者用repartition打散。我一般会在聚合前先按city_id和daterepartition让数据按城市和日期分布均匀然后再做精确聚合。写入 MySQL 时用foreachBatch模式每个批次遍历一次配合事务批量写入比一条条 insert 快一个数量级。val aggDF cleanedDF .withColumn(grid_id, geohash($lon, $lat, 6)) .withColumn(ts_minute, date_trunc(hour, $ts)) .groupBy($city_id, $ts_minute, $grid_id) .agg(countDistinct(bike_id).as(bike_count)) aggDF.writeStream .foreachBatch { (batchDF: DataFrame, batchId: Long) batchDF.write .mode(append) .jdbc( jdbc:mysql://localhost:3306/bike, bike_region_agg, props ) } .outputMode(update) .start() .awaitTermination()geohash在这里把经纬度转成长度为 6 的字符串编码精度大概是 1.2 公里 x 0.6 公里这正好符合运营看热力的需求。如果你要更细的网格可以把长度加到 7但聚合的数据量会呈 4 倍增长MySQL 表的记录数也会爆炸建议先跑一天数据看下网格数再决定精度。4. Spring Boot 服务端把存储变成业务接口4.1 项目搭建与依赖配置三分钟起个 Spring Boot 项目Spring Boot 负责接住上层应用的查询请求它的工程结构不需要搞得很复杂一个bike-storage-service模块就够。你需要引入的依赖主要有五个spring-boot-starter-web、mybatis-spring-boot-starter、mysql-connector-java、hadoop-client和parquet-hadoop。最后两个很多人会漏掉没有它们Spring Boot 就没法直接读 HDFS 上的 Parquet 文件。hadoop-client的版本要跟你 Spark 环境里的 Hadoop 版本对应否则运行时会出现NoSuchMethodError或者ClassNotFoundException。这里踩过一次坑Spark 2.4 配 Hadoop 2.7Spring Boot 里引入了 Hadoop 3.2 的客户端结果连 HDFS 时一直报协议不兼容换了版本才对上。dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.mybatis.spring.boot/groupId artifactIdmybatis-spring-boot-starter/artifactId version2.3.1/version /dependency dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version2.7.7/version /dependency4.2 搭一个接口查区域热力数据运营端最常用的接口是「查某个城市某天各时段的热力数据」。这个查询走的是 MySQL 里的聚合表bike_region_agg数据已经在 Spark 那边算好了Spring Boot 这里只需要做简单的条件查询。用 MyBatis 而不是 JPA 的原因是这种按多条件动态拼 SQL 的场景MyBatis 的Select注解或者 XML 写起来更直接SQL 执行计划也更好控制。RestController RequestMapping(/api/bike/region) public class RegionAggController { Autowired private RegionAggMapper regionAggMapper; GetMapping(/hourly) public ListRegionAggVO hourly( RequestParam Integer cityId, RequestParam String date, RequestParam Integer hour) { return regionAggMapper.selectByCityDateHour(cityId, date, hour); } }控制层逻辑不复杂但注意这里有个性能问题如果运营前端每次都拉全天的热力数据一次会返回几千行记录网络传输和前端渲染都有压力。我一般在查询接口里加一个limit参数限制最大返回条数同时按bike_count倒序排列只返回车辆数最多的前 200 个网格。这样既保证了热力图的完整性又不会把接口拖垮。4.3 读 HDFS 上的 Parquet用 Hadoop 客户端拿到轨迹点文件轨迹回放功能需要读取 HDFS 上 Parquet 格式的原始轨迹数据。这个操作如果直接写在 Controller 里光一个查询就要串行读 100 多个文件每个文件几十 MB接口响应会非常慢。更合理的做法是用异步任务把指定订单的轨迹封装成文件再让前端去下载。但在毕业设计或者内部工具的场景下直接同步读通常也够用。public ListTrajectoryPoint loadTrajectory(String cityId, String date, String orderId) { Configuration conf new Configuration(); conf.set(fs.defaultFS, hdfs://localhost:9000); Path basePath new Path(String.format( /bike_data/trajectory/city_id%s/dt%s, cityId, date)); ListTrajectoryPoint points new ArrayList(); try (ParquetReaderGroup reader ParquetReader.builder(new GroupReadSupport(), basePath).build()) { Group group; while ((group reader.read()) ! null) { String oid group.getString(order_id, 0); if (!orderId.equals(oid)) continue; TrajectoryPoint point new TrajectoryPoint(); point.setLon(group.getDouble(lon, 0)); point.setLat(group.getDouble(lat, 0)); point.setTimestamp(group.getLong(ts, 0)); points.add(point); } } return points; }这段代码里每次读文件都会打开一次 HDFS 连接如果文件多最好对Path做 listFiles 拿到当天所有的 Parquet 文件路径循环读取并合并结果。读取时还有一个隐含问题如果 Spark 作业用的是partitionBy(city_id, dt)目录路径里会有city_idxxx这样的等号结构这是 Hive 风格的分区目录Spring Boot 里手动拼路径时不要搞错层级。5. 避坑与排查把 Spark 和 Spring Boot 数据链路跑通的关键记录5.1 Spark 写出的 Parquet 文件查不到数据现象Spark 作业跑完HDFS 上能看到分区目录文件大小也不是 0但 Spring Boot 查询返回空列表。原因这个情况最常见的原因是 Spark 在写入时没有提交成功或者落盘后文件处于临时状态。检查 HDFS 目录下是不是有.parquet临时文件后缀或者_temporary文件夹如果有说明作业在提交阶段失败了。另一个可能是 Spark 和 Spring Boot 读的 HDFS 路径不一致一个写的是hdfs://localhost:9000另一个读的是本地文件路径file://。解决先把作业的spark.jars和提交方式确认一遍确保写入完成后输出目录下只保留最终的part-*.parquet文件。然后在 Spark 作业结束前手动打印输出路径用hdfs dfs -ls检查分区目录是否存在。Spring Boot 端的读取路径要跟写入时完全一致包括city_idxxx这种分区目录的等号格式。5.2 时区导致的时间字段差了 8 个小时现象聚合结果表里每天的数据少了一部分或者某个小时的数据量明显不对凌晨的数据跑到了前一天。原因Kafka 里的原始轨迹时间戳是 UTC 时间Spark 处理时如果没有把时间戳转成东八区直接按 UTC 日期分区那么北京时间的早上 8 点之前的数据会被划分到前一天。这是共享单车这种国内业务最常见的数据事故不是程序 bug是时区约定问题。解决在 Spark 作业里统一做一次时间转换把 UTC 时间转成Asia/Shanghai时区后再提取日期字段。转换方式是用from_utc_timestamp函数然后基于转换后的时间重算dt分区字段。Spring Boot 查询时也要约定好参数是北京时间不要在服务端再做时区转换否则两次转换会把数据算错。5.3 重复消费导致聚合数据翻倍现象某天的区域热力数据车辆数比实际高一倍部分网格甚至高三四倍。原因Kafka 消费者在作业重启后会从earliest或者上次未提交的 offset 重新消费如果你的作业里没有做幂等写入重复消费一次就会把同一条轨迹点再次聚合。特别是foreachBatch模式下每批数据写 MySQL 时是 append 模式重复批次直接叠加。解决MySQL 侧用INSERT ... ON DUPLICATE KEY UPDATE是第一步让相同主键的记录被覆盖而不是新增。更保险的做法是在 Spark 作业里维护一个 offset 检查点作业重启时从 checkpoint 恢复而不是从 earliest 开始。我一般两种都做checkpoint 管消费位置MySQL upsert 管聚合结果双保险下来数据不会翻倍。5.4 Spark 作业跑着跑着 OOM 挂了现象聚合任务在 shuffle 阶段报OutOfMemoryError作业直接失败。原因共享单车的轨迹数据非常不均匀地铁站附近的点数据密集同一个网格一个小时内可能积累几十万条记录加上countDistinct和geohash这两个操作都需要大量内存。默认的 Spark 内存配置是按通用场景设置的没有为这种高度倾斜的数据做调整。解决调两个参数spark.sql.shuffle.partitions从默认的 200 调到实际 Executor 核心数的 2 到 3 倍spark.executor.memory如果机器内存够就给到 4G 以上。如果数据倾斜问题依然存在在聚合前先用repartition把热点 key 打散。5.5 Spring Boot 启动时跟 Spark 抢内存现象Spring Boot 应用启动后用一段时间就频繁 Full GC接口响应越来越慢。原因有人会在 Spring Boot 进程里直接调用 SparkSession 做查询两个框架都吃内存而且 Spark 的local[*]模式会把本机所有核心拿来跑完全不考虑 Spring Boot 还有一堆连接池要跑。这是架构上的错误不是配置能解决的。解决Spring Boot 里不要初始化 SparkSessionSpark 作业独立提交Spring Boot 只负责查结果数据要用 HDFS 上的轨迹明细时用 Hadoop 客户端轻量读取。切割离线计算和在线服务的边界是这类系统最容易忽略的设计问题。6. 怎么验证这套存储系统真的没问题三个能落地的检查方法6.1 数据对账Spark 结果和原始日志的条数对比系统上线前至少要对一天的数据做一次全量对账。做法很简单用 Spark 读原始日志统计总条数再读清洗后的 Parquet 统计落库条数两个数对比差异率超过 1% 就要查清洗规则是不是太严格了。下面这段脚本可以在 Spark Shell 里直接跑val rawCount spark.read.json(/data/bike_raw).count() val cleanCount spark.read.parquet(/bike_data/trajectory).count() println(sraw$rawCount clean$cleanCount lost${rawCount - cleanCount})正常清洗规则下丢失率应该在 0.5% 以下。如果丢得太多多半是你的经纬度范围过滤写错了把国内城市的数据过滤掉不少。这个对账脚本建议做成每天定时任务自动输出对比结果数据质量出问题能第一时间发现。6.2 查询性能验证热点查询必须控制在 200ms 内对 Spring Boot 接口做简单的压测关注两个数字P99 延迟和吞吐量。区域热力查询走 MySQL 索引200ms 以内算正常轨迹回放查询要读 HDFS 文件如果超过 1 秒说明文件数量太多或者单文件太大考虑用更小的分区粒度重新落数据。验证时用 JMeter 或者直接写 Python 脚本发起并发请求观察数据库连接池是否被打满。6.3 进阶改进把 Spark 作业升级成 Structured Streaming 实时链路当批处理链路稳定之后你可以把 Spark 作业从readStream换成真正的持续处理模式把触发间隔从 10 分钟缩短到 1 分钟。这样 MySQL 聚合表的数据从日级变成近实时运营看到的热力数据就接近当前时刻了。这个升级不需要改表结构只需要调整 trigger 和检查点位置但要注意 MySQL 的写入频率会提升 10 倍连接池要相应加大否则会成为新的瓶颈。做这套系统的过程中我最大的一个教训是不要一开始就把架构设计得太复杂先让 Spark 和 Spring Boot 各管一段跑通一条链路再逐步加实时性和性能优化。很多所谓的大数据系统最终死在不必要的复杂度上而不是死在数据量上。希望这套「Spark 算、Spring Boot 查」的思路能给你一条可靠的技术路线帮你在做共享单车数据存储系统时少走几段弯路。本文还有配套的精品资源点击获取
返回列表