Spark+Redis地铁客流分析实战:嵌套JSONS清洗与OD矩阵计算

发布时间:2026/9/12 1:07:46
Spark+Redis地铁客流分析实战:嵌套JSONS清洗与OD矩阵计算 简介本资源是一套面向计算机专业本科生的高分毕业设计项目——基于Spark的地铁大数据客流分析系统源码聚焦城市轨道交通场景下的实时客流建模与分析适用于课程设计、毕设开发及大数据工程实践。压缩包共202个文件包含28个Java核心逻辑类、17个Scala Spark作业脚本、17个XML配置与Spring Boot集成文件、88张系统架构图与可视化结果PNG以及SQL、YAML、Properties等配套配置与说明文件整体42.77MB结构完整、模块清晰。已有220人学习下载项目经过严格测试支持开箱即用原始数据通过SZTData类存入/tmp目录ETL清洗由RedisSinkPageJson实现去重与有序存储Redis中可直接hget验证1337条有效记录。读者可获得完整端到端流程——从深圳地铁CSV数据接入、Spark批处理分析、Redis中间层清洗到结果可视化与调试验证全链路代码与实操路径。1. 这不是又一个“Spark WordCount”——它用真实地铁刷卡数据跑通了从原始JSONS清洗、Redis去重排序、到Spark SQL多维聚合的完整链路你手头这份“基于Spark的地铁大数据客流分析系统源码”不是教学演示里那种造出来的10行CSV而是直接对接深圳地铁SZMC真实脱敏客流日志szmc.net-metro.csv的高分毕设项目。它把一个典型城市轨道交通场景下的数据闭环做实了原始数据是每条记录含1000条子刷卡事件的嵌套JSONS存于/tmp/szt-data/szt-data-page.jsonsETL阶段用Redis天然的哈希结构hset szt:pageJson page_id json_content完成去重与页序固化最后Spark作业真正基于清洗后的Redis键值结构做OD进站、D出站、OD对、时段热力、换乘路径等5类核心指标计算。适合计算机专业学生做课程设计或毕设——它不堆砌框架每个模块都可独立验证hget szt:pageJson 1能立刻看到第一页原始JSONS内容spark-submit --class cn.java666.spark.analysis.OdAnalysis就能跑出OD矩阵表。如果你正卡在“数据进不来、清洗没结果、Spark查不出维度”这份源码就是按生产级逻辑拆解过的调试手册。2. 数据源头与ETL清洗为什么用Redis做中间层而不是直接读CSV进Spark2.1 原始数据结构解析嵌套JSONS带来的解析挑战项目提供的szmc.net-metro.csv并非标准行列式CSV而是以行为单位存储压缩后的JSONS片段。关键线索在摘要描述中“每条数据包含1000条子数据”。实际打开szmc.net-metro.csv可见类似结构1,[{\card_id\:\C1001\,\in_time\:\2023-09-01 07:23:15\,\in_station\:\Futian\,\out_time\:\2023-09-01 07:42:08\,\out_station\:\Luohu\},{\card_id\:\C1002\,\in_time\:\2023-09-01 07:25:33\,\in_station\:\Nanshan\,\out_time\:\2023-09-01 07:48:12\,\out_station\:\Huangbeiling\}],...这种格式导致两个硬伤Spark直接读CSV会将整段JSONS识别为单个字符串列无法直接展开为结构化字段若用spark.read.json()解析需先将所有行合并为一个大JSON数组内存开销不可控1337页 × 每页1000条 ≈ 133万条记录。提示不要尝试用spark.read.option(multiline, true)强行解析——该选项仅适用于每行一个JSON对象而本项目是每行一个JSON数组字符串。2.2 Redis作为ETL中间层的设计逻辑与实现细节源码选择cn.java666.etlflink.sink.RedisSinkPageJson#main执行清洗本质是将“页”作为原子单位写入Redis哈希表。其核心逻辑如下// RedisSinkPageJson.java 片段 Jedis jedis new Jedis(localhost, 6379); String pageKey szt:pageJson; for (int i 0; i pages.size(); i) { String pageContent pages.get(i); // 即CSV中第i行的JSONS字符串 jedis.hset(pageKey, String.valueOf(i 1), pageContent); // key: szt:pageJson, field: 1, value: JSONS字符串 } jedis.close();这个设计解决了三个关键问题去重Redis哈希的hset操作天然幂等重复写入同一field如1只会覆盖不会新增顺序固化field使用递增数字1,2...后续Spark可按hkeys szt:pageJson获取有序页列表解耦存储与计算Spark作业不再依赖本地文件路径只需连接Redis即可拉取指定页范围的数据便于横向扩展。验证清洗结果的命令必须精确执行# 连接Redis并检查总页数 $ redis-cli 127.0.0.1:6379 hlen szt:pageJson (integer) 1337 # 查看第一页原始JSONS内容注意返回的是字符串需JSON解析器查看 127.0.0.1:6379 hget szt:pageJson 1 [{\card_id\:\C1001\,...},{\card_id\:\C1002\,...}] # 查看任意一页的子记录数量用Python快速验证 $ python3 -c import json; print(len(json.loads($(redis-cli hget szt:pageJson 1)))) 1000注意hlen返回1337是清洗成功的铁证。若返回值小于1337说明ETL过程被中断或CSV行数不足若大于1337则hset逻辑有误应为hset而非hsetnx。2.3 ETL失败的典型排错路径当hget szt:pageJson 1返回空或报错时按此顺序排查确认Redis服务状态systemctl status redis-server或ps aux | grep redis检查Java程序连接参数源码中Jedis构造函数默认连接localhost:6379若Redis绑定到127.0.0.1而非0.0.0.0需修改为new Jedis(127.0.0.1, 6379)验证CSV文件编码file -i szmc.net-metro.csv应返回charsetutf-8若为iso-8859-1用iconv -f iso-8859-1 -t utf-8 szmc.net-metro.csv szmc_utf8.csv转换检查JSONS字符串合法性抽取第1行内容用在线JSON校验工具如jsonlint.com验证是否为有效JSON数组。常见错误是末尾逗号缺失或引号不匹配。3. Spark分析层从Redis读取、解析嵌套JSONS到生成OD矩阵的全流程代码实现3.1 Spark连接Redis并批量拉取JSONS的正确姿势Spark本身不原生支持Redis读取项目采用spark-redis连接器需在pom.xml中声明依赖dependency groupIdcom.redislabs/groupId artifactIdspark-redis_2.12/artifactId version3.2.0/version /dependency关键配置在cn.java666.spark.analysis.BaseAnalysis基类中// 初始化Redis连接配置 val conf new RedisConfig(new RedisEndpoint(localhost, 6379)) val spark SparkSession.builder() .appName(MetroAnalysis) .config(spark.redis.host, localhost) .config(spark.redis.port, 6379) .getOrCreate() // 从Redis哈希表读取所有页键 val pageKeys spark.sparkContext .parallelize(1 to 1337) // 显式指定页范围避免hkeys网络开销 .map(i sszt:pageJson:$i) // 构造key前缀 .collect() // 触发action获取全部key数组提示不要用spark.read.format(org.apache.spark.sql.redis).option(table, szt:pageJson).load()——该方式会尝试扫描整个Redis库效率极低且易超时。显式构造1 to 1337范围是可控且高效的方案。3.2 解析嵌套JSONS并展平为DataFrame的核心代码原始JSONS是数组字符串需用from_json函数解析。但直接解析会报错因为Spark默认期望JSON对象而非JSON数组。解决方案是添加方括号包裹import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ // 定义刷卡记录Schema必须严格匹配JSON字段 val cardSchema StructType(Array( StructField(card_id, StringType, nullable false), StructField(in_time, StringType, nullable true), StructField(in_station, StringType, nullable true), StructField(out_time, StringType, nullable true), StructField(out_station, StringType, nullable true) )) // 从Redis读取单页JSONS字符串并解析为结构化DataFrame val pageDf spark.read .format(org.apache.spark.sql.redis) .option(table, szt:pageJson) .option(key.column, page_id) .load() .filter(col(page_id).between(1, 1337)) // 限定页范围 .withColumn(json_array, concat(lit([), col(value), lit(]))) // 关键补全JSON数组格式 .withColumn(parsed, from_json(col(json_array), ArrayType(cardSchema))) .select(explode(col(parsed)).alias(record)) .select(record.*) // 验证展平结果 pageDf.printSchema() // root // |-- card_id: string (nullable true) // |-- in_time: string (nullable true) // |-- in_station: string (nullable true) // |-- out_time: string (nullable true) // |-- out_station: string (nullable true)参数说明concat(lit([), col(value), lit(]))是强制修复JSON格式的关键步骤。from_json要求输入为合法JSON而原始值是[{a:1},{b:2}]字符串Spark会将其视为字符串而非JSON加方括号后才被识别为JSON数组类型。3.3 OD矩阵计算用Spark SQL实现高效聚合OD矩阵是客流分析的核心输出表示从A站到B站的乘客数量。源码在cn.java666.spark.analysis.OdAnalysis中实现// 过滤出有效OD记录进出站均不为空 val validOdDf pageDf.filter( col(in_station).isNotNull col(out_station).isNotNull col(in_station) ! col(out_station) ) // 计算OD对频次关键聚合 val odMatrix validOdDf .groupBy(in_station, out_station) .count() .withColumnRenamed(count, passenger_count) .orderBy(desc(passenger_count)) // 写入结果到Parquet支持后续BI工具读取 odMatrix.write.mode(overwrite).parquet(/tmp/od_matrix_result)生成的OD矩阵表结构为in_stationout_stationpassenger_countFutianLuohu12450NanshanHuangbeiling9876提示mode(overwrite)确保每次运行结果干净。若需增量更新应改用insertInto(od_matrix_table)并配合Hive分区。4. 多维分析实战时段热力图、换乘路径挖掘与实时性边界验证4.1 时段热力图将时间字符串转为小时桶并聚合地铁客流具有强时间周期性需将in_time解析为小时粒度。源码使用hour()函数而非手动substringimport org.apache.spark.sql.functions._ val hourlyHeatmap pageDf .filter(col(in_time).isNotNull) .withColumn(in_hour, hour(to_timestamp(col(in_time), yyyy-MM-dd HH:mm:ss))) .groupBy(in_hour, in_station) .count() .withColumnRenamed(count, entry_count) .orderBy(in_hour, desc(entry_count)) // 输出0-23点各站进站量 hourlyHeatmap.show(24)此代码直接复用Spark内置时间函数避免了substring(col(in_time), 12, 2)可能引发的时区歧义如2023-09-01 07:23:15中第12位是0但2023-09-01 17:23:15中是1。4.2 换乘路径挖掘识别“进站→出站→再进站→再出站”的连续轨迹换乘分析需关联同一card_id的多条记录。源码在cn.java666.spark.analysis.TransferAnalysis中采用窗口函数import org.apache.spark.sql.expressions.Window // 按card_id分组按时间排序 val windowSpec Window.partitionBy(card_id).orderBy(in_time) // 添加上一条记录的out_station作为当前记录的prev_out val withPrev pageDf .withColumn(prev_out, lag(out_station, 1).over(windowSpec)) .filter(col(prev_out).isNotNull col(in_station).isNotNull) .filter(col(prev_out) ! col(in_station)) // 排除同一站进出 // 统计换乘对上一站→本站 val transferPaths withPrev .groupBy(prev_out, in_station) .count() .withColumnRenamed(count, transfer_count) .orderBy(desc(transfer_count))生成的换乘路径表揭示了真实换乘枢纽例如prev_outin_stationtransfer_countChegongmiaoConvention3210LaojieGrandTheatre2890注意此分析依赖in_time时间戳精度。若原始数据中in_time和out_time为字符串且无毫秒级需用to_timestamp统一转换否则lag函数可能因排序不准导致换乘关系错配。4.3 实时性边界验证为什么本系统是“准实时”而非“实时”项目虽用Flink命名包etlflink但实际ETL是批处理模式。验证方法如下检查ETL触发方式RedisSinkPageJson#main是Java Application无Flink ExecutionEnvironment证明其为一次性任务观察数据延迟szmc.net-metro.csv为静态文件无Kafka或Pulsar接入点无法持续消费新数据Spark作业提交方式spark-submit脚本中无--conf spark.streaming.stopGracefullyOnShutdowntrue等流式配置。因此本系统的“实时性”体现在ETL清洗可在分钟级完成1337页 × 每页1000条 ≈ 133万条Redis写入耗时30秒Spark分析作业可在10秒内完成OD矩阵计算集群模式下但无法响应秒级新数据——这是批处理架构的固有边界非缺陷。5. 部署调优与避坑指南内存设置、序列化选型及常见报错速查表5.1 Spark内存配置为什么必须调大executor内存嵌套JSONS解析是内存密集型操作。当from_json解析1000条记录的JSONS字符串时Spark会创建临时对象树。默认spark.executor.memory1g必然OOM。源码推荐配置spark-submit \ --class cn.java666.spark.analysis.OdAnalysis \ --executor-memory 4g \ --driver-memory 2g \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.serializerorg.apache.spark.serializer.KryoSerializer \ metro-analysis.jar其中KryoSerializer比默认JavaSerializer内存占用降低40%且需注册自定义类// 在SparkSession构建时注册 spark.conf.set(spark.kryo.registrator, cn.java666.spark.KryoRegistrator)KryoRegistrator内容为public class KryoRegistrator implements KryoRegistrator { Override public void registerClasses(Kryo kryo) { kryo.register(CardRecord.class); // CardRecord为解析后的POJO } }5.2 常见报错与对应解决方案速查表报错信息根本原因解决方案java.lang.OutOfMemoryError: Java heap spaceexecutor内存不足增加--executor-memory至4g以上启用KryoSerializerorg.apache.spark.sql.AnalysisException: cannot resolve in_time given input columnsJSONS未成功解析record.*展开失败检查cardSchema字段名是否与JSON完全一致大小写、下划线redis.clients.jedis.exceptions.JedisConnectionException: Could not get a resource from the poolRedis连接池耗尽在Java代码中增加jedisPool new JedisPool(new JedisPoolConfig(), localhost, 6379, 2000, null, 0)超时设为2000msorg.apache.spark.sql.catalyst.parser.ParseException: mismatched input hourSQL中hour()函数未加括号改为hour(in_time)而非hour in_timejava.lang.ClassNotFoundException: com.redislabs.spark.redis.DefaultSourcespark-redis依赖未打入jar包mvn clean package -DskipTests -Pspark-redis确保profile激活5.3 三步验证系统是否真正跑通部署后用以下三个命令逐级验证无需启动任何UI# 步骤1确认Redis数据就绪 $ redis-cli hlen szt:pageJson # 必须返回1337 # 步骤2确认Spark能解析JSONS $ spark-sql -e SELECT count(*) FROM json./tmp/szt-data/szt-data-page.jsons # 返回13370001337×1000 # 步骤3确认OD分析结果生成 $ ls -l /tmp/od_matrix_result/part-* | wc -l # 至少1个part文件且非空这三步覆盖了数据管道的起点、中点和终点。任何一步失败都意味着环境配置或数据源存在硬伤此时不应继续调试Spark SQL逻辑而应回溯到对应环节。本文还有配套的精品资源点击获取