多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

Spark与Spring Boot构建共享单车数据存储系统:架构与实战

Spark与Spring Boot构建共享单车数据存储系统:架构与实战 简介面向毕业设计或课程实践的SpringBoot与Spark共享单车数据存储系统全套资料适合Java方向及大数据相关专业学生。内容涵盖完整前后端源码、SQL脚本、部署运行脚本与配套论文文档可支撑从环境搭建、功能调试到论文撰写的全过程。资源共374个文件以Java后端、Vue前端、JavaScript脚本、SVG图标及配置文件为主辅以数据库初始化文件与说明文档压缩包约17.8MB结构清晰便于按模块查阅。已有51人学习下载可用作共享单车数据采集、存储与分析场景的参考实现帮助理解分布式存储、Spark数据处理及RESTful API设计等知识点。作者支持私信交流技术问题下载后建议先查看README或论文文件仅限学习交流使用。1. 共享单车数据存储系统为什么 Spark 和 Spring Boot 的搭配是当下性价比最高的方案共享单车一秒内产生的订单和轨迹数据就有上千条一天下来就是几千万的量级。用 MySQL 硬扛写入先扛不住查询端一上来就是慢查询。用纯 Hadoop 生态写入是解决了但 Spring Boot 服务层每次查数据都要写 MapReduce开发效率太低。这个标题的价值在于它把存储系统拆成了两层底层用 Spark 做清洗和落盘上层用 Spring Boot 做查询接口。既不放弃大数据的吞吐能力也不牺牲 Java 工程化的开发效率。省市级共享单车监管平台、高校毕设项目、企业内部的出行数据分析中台都适用这条技术路径。今天这篇直接把架构选型、表设计、核心代码和踩坑记录全部分享出来。2. 先想清楚整条存储链路再写代码数据从采集到可查询的四个环节2.1 为什么选 Spark 而不是 Flink批处理在存储系统里的优势这个系统的核心定位是存储不是实时告警。存储系统的第一诉求是数据完整性和可回溯性不是毫秒级延迟。Spark 的 Structured Streaming 用微批处理模型每一批数据都有明确的提交边界出错了可以从 checkpoint 恢复重跑。Flink 在流处理上的延迟确实更低但状态管理和 checkpoint 的调优复杂度更高对于论文和中小型团队来说维护成本偏大。Spark 另一个优势是它可以一套代码同时跑批处理和准实时处理。白天用 Structured Streaming 消费 Kafka 里的实时订单流晚间用同一个 SparkSession 跑全量补数任务代码复用率高。Spark SQL 的 DataFrame API 做列裁剪和谓词下推很成熟对数据清洗来说是天然的优势。Spark 集群我一般建议用 Standalone 模式而不是 YARN原因是共享单车存储系统的数据规模通常在 TB 级以下不值得为一个存储系统单独维护一套 Hadoop YARN。Standalone 模式部署简单三台机器就能撑起千万级日订单量的清洗任务。2.2 存储层选型Kafka、HDFS、HBase、Redis 各管一段共享单车数据存储系统最忌讳的就是一张表存天下。我的分层方案是四层层级组件职责数据留存接入缓冲Kafka接收订单和轨迹原始数据削峰填谷7 天离线存储HDFS存储清洗后的 Parquet 文件做冷数据备份永久在线查询HBase支撑按订单号、时间范围的实时查询90 天热点加速Redis缓存最近一小时高频查询结果1 小时Kafka 的 Topic 按业务类型拆分order-topic存订单数据gps-topic存车辆轨迹数据。HDFS 上按日期分区存 Parquet 格式文件压缩比高配合 Spark SQL 查询效率好。HBase 的 RowKey 设计是这套系统的关键后面会单独讲。2.3 订单表和轨迹表怎么设计HBase RowKey 与 HDFS 分区规则订单表的 RowKey 我采用bikeId reverse(orderTime)的设计。反转时间戳是为了让同一辆车的最近订单排在前面方便查询某辆车最近的使用记录。轨迹表用orderId timestamp拼接保证一次查询能连续扫描出完整轨迹。HDFS 上的 Parquet 分区按date/hour两层分区。为什么按小时再分一层因为夜间是共享单车的使用低谷如果不按小时分区每天凌晨的 Spark 任务要扫描全天数据白白浪费大量资源。CREATE TABLE IF NOT EXISTS ods_bike_order ( order_id STRING, bike_id STRING, user_id STRING, start_lng DOUBLE, start_lat DOUBLE, end_lng DOUBLE, end_lat DOUBLE, start_time TIMESTAMP, end_time TIMESTAMP, duration_sec INT, distance_km DOUBLE ) USING parquet PARTITIONED BY (date STRING, hour STRING);分区字段放在表结构最后Spark SQL 写入时会自动按分区目录组织文件。start_time 和 end_time 用 TIMESTAMP 类型而不是 STRING避免后续做时长计算时每次都要 cast。3. 用 Spring Boot 把 Spark 存储系统跑起来从依赖到可查询的关键代码3.1 项目结构与 Maven 依赖一套代码同时跑任务和接口我习惯把项目拆成engine和web两个 module。engine 里放 Spark 任务代码打成 fat jar 提交到集群web 里放 Spring Boot 接口代码处理查询请求。为什么要拆因为 Spark 的依赖和 Spring Boot 的依赖存在冲突尤其commons-logging和log4j这两个库两边版本不一致时放在同一个工程里会出现类加载混乱。dependency groupIdorg.apache.spark/groupId artifactIdspark-sql_2.12/artifactId version3.3.2/version scopeprovided/scope /dependency dependency groupIdorg.apache.spark/groupId artifactIdspark-streaming-kafka-0-10_2.12/artifactId version3.3.2/version /dependency dependency groupIdorg.apache.hbase/groupId artifactIdhbase-client/artifactId version2.4.17/version /dependencySpark 相关依赖加scopeprovided/scope提交任务时用集群自带的 Spark 环境避免把 Spark 的类打进 fat jar 里。HBase 客户端单独放因为 Spring Boot 的 web 模块也需要连 HBase 做在线查询。3.2 SparkSession 初始化与参数配置三个必调的参数SparkSession 的初始化放在工具类里统一管理。这里有三组参数我每次都会调分别是spark.sql.shuffle.partitions、spark.sql.adaptive.enabled和spark.serializer。val spark SparkSession.builder() .appName(SharedBikeStorageETL) .master(spark://node01:7077) .config(spark.sql.shuffle.partitions, 48) .config(spark.sql.adaptive.enabled, true) .config(spark.serializer, org.apache.spark.serializer.KryoSerializer) .config(spark.sql.parquet.compression.codec, snappy) .config(spark.streaming.kafka.consumer.cache.enabled, false) .enableHiveSupport() .getOrCreate()spark.sql.shuffle.partitions默认值是 200但共享单车清洗任务里 shuffle 的数据量不大200 个分区会产生大量空 task。48 这个值对应 3 台 worker 节点每台 16 核。spark.sql.adaptive.enabled打开后Spark 会在运行时合并小 shuffle 分区进一步减少空任务。spark.serializer换成 Kryo 后RDD 在 shuffle 和落盘时的序列化开销能降低 30% 以上。3.3 消费 Kafka 订单流Structured Streaming 写入 HDFS 的完整链路共享单车的原始订单 JSON 直接从 Kafka 消费这里我用 Structured Streaming 的readStreamAPI 读取然后做清洗后写入 HDFS。val kafkaParams Map( kafka.bootstrap.servers - node01:9092,node02:9092, subscribe - order-topic, startingOffsets - earliest, maxOffsetsPerTrigger - 100000, failOnDataLoss - false ) val rawStream spark.readStream .format(kafka) .options(kafkaParams) .load() .selectExpr(CAST(value AS STRING) as json_str) val parsedStream rawStream .select(from_json(col(json_str), orderSchema).alias(data)) .select(data.*) .filter(col(order_id).isNotNull) .filter(col(start_time).isNotNull)maxOffsetsPerTrigger设置为 100000这是控制单批次处理量的关键参数。设置太小处理速度跟不上数据产生速度造成 Kafka 消息积压设置太大单批次处理时间过长微批延迟变高。另外务必要设置failOnDataLoss为false否则 Kafka 日志过期被清理后任务直接崩溃。3.4 清洗逻辑去重、时区修正、距离字段补全共享单车真实数据里常见的脏数据有三种重复订单、时区错乱、坐标异常。重复订单在 Kafka 因为消费者重平衡可能产生重复投递需要按order_id去重。val deduped parsedStream.dropDuplicates(order_id) val cleaned deduped .withColumn(start_time, from_utc_timestamp(to_timestamp(col(start_time), yyyy-MM-ddTHH:mm:ssZ), Asia/Shanghai)) .withColumn(end_time, from_utc_timestamp(to_timestamp(col(end_time), yyyy-MM-ddTHH:mm:ssZ), Asia/Shanghai)) .withColumn(duration_sec, unix_timestamp(col(end_time)) - unix_timestamp(col(start_time))) .withColumn(distance_km, round(haversine(col(start_lng), col(start_lat), col(end_lng), col(end_lat)), 2)) .filter(col(start_lat).between(22.4, 41.2)) .filter(col(start_lng).between(101.2, 121.8))dropDuplicates需要指定order_id字段Spark 会用这个字段做全分区 shuffle 去重数据量大时比较耗资源。如果确认 Kafka 的enable.idempotence已经打开可以去掉去重步骤。时区修正这里是个大坑原始数据里的时间戳是 UTC 的 ISO8601 格式如果不转北京时间后续按小时分区时所有订单的时间都会错 8 个小时。3.5 Spring Boot 查询服务从 HBase 读取订单数据返回前端数据落到 HBase 之后Spring Boot 的服务层负责读取并返回给前端。这里我用 HBase 的 Java 客户端 API 直接查询不用 Phoenix因为查询模式固定就是按键查不需要复杂 join。RestController RequestMapping(/api/order) public class OrderController { Autowired private HBaseQueryService hBaseQueryService; GetMapping(/{orderId}) public ResultOrderVO getOrder(PathVariable String orderId) { OrderVO order hBaseQueryService.findByOrderId(orderId); return Result.success(order); } GetMapping(/bike/{bikeId}/recent) public ResultListOrderVO getRecentOrders( PathVariable String bikeId, RequestParam(defaultValue 10) int limit) { ListOrderVO orders hBaseQueryService.findRecentByBikeId(bikeId, limit); return Result.success(orders); } }HBase 查询第一步先用 Redis 查缓存命中直接返回没命中再去 HBase 查查完回填 Redis 并设置过期时间。为什么缓存只存一小时因为共享单车订单数据查询有很明显的时间聚集效应比如早高峰时段用户密集查询某几个站点的车辆订单一小时后热点转移缓存自动失效。4. 落地这套共享单车存储系统的五个翻车场景与排查清单4.1 Spark 任务 OOM分区数没算对Executor 直接挂掉现象任务运行到 stage 2 时 Executor 报java.lang.OutOfMemoryError重试三次后整个 application 失败。原因HDFS 上的 Parquet 文件在扫描时默认按文件分片每个分片对应一个 task。如果每天生成的文件数超过 1000 个Spark 会为每个文件开一个 taskExecutor 内存被 task 并发数冲垮。这是共享单车数据按小时分区后最常见的翻车点。解决在读取时设置spark.sql.files.maxPartitionBytes为 256MB强制 Spark 合并小文件为大分区。同时调大 Executor 内存spark.executor.memory设置为 4gspark.executor.cores设置为 4每个 Executor 上并发 task 数控制在 4 个以内。4.2 HBase RowKey 设计不合理导致的热点问题现象写入 HBase 时 region server 频繁报Blocking wait for lock部分节点负载飙升其他节点空闲。原因订单表 RowKey 用bikeId reverse(orderTime)后如果按 bikeId 前缀写入某一辆热门车辆的订单会持续写入同一个 region形成写热点。解决在 RowKey 前拼接一个 2 位的散列前缀用bikeId.hashCode % 97计算得出。这样写入会均匀分布到 97 个 region 中。查询时先算出散列值再拼接完整 RowKey查询逻辑不受影响。这是 HBase 表设计中最容易被忽略的细节等线上出问题再改 RowKey 就需要迁移全表数据。4.3 Kryo 序列化报错自定义类没注册任务启动就失败现象任务提交后 executor 日志报ClassNotFoundException指向自定义的订单结构体。原因Spark 默认的 Java 序列化可以处理任意类但效率低换成 Kryo 后需要手动注册自定义类型否则 Kryo 不知道如何处理。解决在构建 SparkConf 时用registerKryoClasses显式注册。val sparkConf new SparkConf() .setAppName(SharedBikeStorageETL) .set(spark.serializer, org.apache.spark.serializer.KryoSerializer) .registerKryoClasses(Array( classOf[OrderRecord], classOf[GpsRecord], classOf[Array[OrderRecord]] ))还有一个隐藏的坑dropDuplicates操作内部会产生无界状态存储如果没开spark.sql.streaming.statefulOperator.checkConfirmation或没有设置 state TTL状态体积会无限膨胀。建议在流任务里给去重加上时间窗口限制只保留最近 24 小时的状态。4.4 小文件问题凌晨跑完批任务HDFS 上多出上万个文件现象HDFS NameNode 内存告警查询 Parquet 表时 Spark 列出文件列表耗时超过 10 秒。原因Structured Streaming 每个微批处理完数据都会写新文件如果每 5 秒一个批次一天下来会产生 17280 个文件。还没算上按小时分区的目录开销。解决在流任务里加一个coalesce或repartition控制输出分区数强制每个批次的输出文件数不超过当前目录下文件数的预期值。另外我通常在凌晨加一个合并小文件的批任务用OPTIMIZE语法合并当天的小 Parquet 文件。INSERT OVERWRITE TABLE ods_bike_order PARTITION (date2024-11-20, hour08) SELECT * FROM ods_bike_order WHERE date2024-11-20 AND hour08 DISTRIBUTE BY order_id;DISTRIBUTE BY在这里做了一次 shuffle结果写入时每个 reducer 只写 1 个文件把原来上百个小文件压缩成可控数量。4.5 查询接口时快时慢连接池没预热第一次查询凉凉现象Spring Boot 服务重启后前几个查询接口响应时间长达 5 秒后续恢复为 50ms。原因HBase 连接池在第一次请求时初始化同时 HBase 的 region 定位需要扫描 meta 表第一次访问没有缓存。解决在 Spring Boot 启动类里用ApplicationRunner预热连接池项目启动时主动执行一次 HBase 表的exists调用和一次空查询。Component public class HBasePreheater implements ApplicationRunner { Override public void run(ApplicationArguments args) { for (String table : tableConfig.getAll()) { hBaseQueryService.exists(table); } log.info(HBase connection pool preheated); } }5. 数据质量校验与性能验证上线前必须做透的两件事5.1 数据抽样校验拿原始 JSON 对比 HDFS 落盘表存储系统的核心是数据没丢、没串、没重复。我的校验方法是每天跑完批任务后从 Kafka 的原始 Topic 里随机取 1 万条数据和 HDFS 上的ods_bike_order表做对比。校验用 Spark SQL 写一个简单的比对任务SELECT count(*) AS total_cnt, count(DISTINCT order_id) AS distinct_cnt, count(IF(start_time 2024-01-01 OR start_time 2025-01-01, 1, NULL)) AS time_outlier_cnt, count(IF(start_lng 100 OR start_lng 122, 1, NULL)) AS lng_outlier_cnt FROM ods_bike_order WHERE date 2024-11-20;total_cnt和distinct_cnt的差值如果在 0.1% 以内说明去重逻辑生效。time_outlier_cnt大于 0 说明时区转换还有边界情况没处理。lng_outlier_cnt大于 0 说明部分设备上报的定位数据本身有问题需要在上游过滤。5.2 性能基准从 Kafka 到可查询的端到端耗时我会做两组性能测试来评估系统是否达标。第一组是写入链路向 Kafka 灌入 100 万条模拟订单数据记录从数据进入 Kafka 到 HDFS 落盘完成的时间。目标是不超过 5 分钟。第二组是查询链路通过 Spring Boot 接口并发 100 个线程查询随机订单号观察 HBase 读吞吐和响应时间分布。如果读链路 qps 上不去优先检查 region server 的堆内存分配。HBase 的hbase.regionserver.global.memstore.size默认是堆的 40%如果写压力大读性能会被 flush 阻塞拖垮。我会把这个参数调到 30%把释放出来的内存留给 block cache设置hfile.block.cache.size为 35%。最后一件事监控不能省。我在 Prometheus 里用spark_streaming_processing_delay和kafka_consumer_lag这两个指标盯住链路健康度。处理延时一旦超过 2 分钟或者 lag 持续增长直接告警到手机。因为存储系统的故障不会立刻暴露往往是三天后查数据才发现少了某一天的记录那时候想用 Kafka 里的原始数据做补偿可能已经被日志清理策略删掉了。这是我吃过亏后总结下来的习惯希望帮到你。本文还有配套的精品资源点击获取
返回列表