基于Spark的地铁客流分析:从AFC数据到可视化大屏实战
2026/9/14 17:28:25 网站建设 项目流程

简介:一份基于Spark的地铁大数据客流分析系统源码,面向计算机科学与大数据相关专业的毕业设计学生,也适合想快速上手Spark离线处理的开发者。项目覆盖数据采集、清洗、分析到可视化展示的完整链路,可帮助理解Spark RDD与DataFrame等常用算子在地铁刷卡数据上的实际应用。压缩包共202个文件,约42.77MB,其中以Java与Scala源码为核心,配以yaml、xml、properties等环境配置,以及png图片和sql脚本,便于还原项目结构与数据库表设计,另有部分md说明文档辅助部署。目前已有327人浏览学习,项目已经过本地编译验证,下载后按配套说明配置环境即可运行,整体难度适中,能有效支撑毕业设计选题、代码参考与答辩讲解。

1. 基于Spark的地铁大数据客流分析,难在把海量刷卡记录变成调度指令

地铁闸机的AFC系统一天会产生上千万条刷卡记录,全网OD组合以十万计。客流分析要回答的其实只有三个问题:某站在某个时段进出多少人、哪些区间最拥挤、换乘压力集中在哪。这些答案直接决定行车交路怎么排、限流措施几点启动。用Spark来做,不是因为单机跑不动,而是它能把清洗、OD聚合、断面推算放进同一条流水线,从单机到集群迁移成本几乎为零。下面按架构选型、集群搭建、指标计算、落库可视化、答辩验证的顺序推进,适合正在选型、或者代码写到一半不知道怎么组织模块和源码结构的人。

2. 客流分析系统的选型依据与Spark集群搭建

2.1 为什么选Spark而不是MapReduce或Flink

选型先看数据形态。AFC原始数据是日增上千万行的结构化流水,业务上要的是"每天定时出结果",属于典型的批处理场景。MapReduce能算,但映射到OD矩阵和断面客流时,每加一个指标就要重写一套Java MR程序,调试一轮的代价足够写完三个DataFrame任务。Flink适合秒级延迟的实时限流预警,但状态管理、watermark、checkpoint这些概念,在毕设周期里会吃掉大量本应用于结果分析的时间,而且答辩时很难一句话讲清状态后端。

折中下来Spark是更稳妥的选项。它可以以Structured Streaming的形式做分钟级微批,也可以用DataFrame批量算日结指标,一套API覆盖两种模式。下面这张表是我在选型时会横向比的内容,也是答辩"为什么不用XX"的直接依据。

计算框架计算模型指标开发成本集群运维适合的客流场景
MapReduce纯批量高,每个指标都要写MR周级或月级离线报表
Spark微批+批量低,DataFrame或SQL均可日级OD矩阵、断面客流、站点热度
Flink真流式中,状态管理复杂秒级实时限流预警

还有一个容易被忽略的点:Spark的DataFrame在Catalyst优化器下会自动做谓词下推和列裁剪,同样的聚合逻辑比手写RDD算子的代码短三分之一,执行计划更稳定。毕设源码的可读性直接影响评分,从RDD迁移到DataFrame这一步投入产出比很高。

提示:如果导师更看重实时性,用Spark Structured Streaming可以边跑边答"为什么不用Flink",但核心指标仍然建议批量计算后落库,实时链路只做增量刷新。

2.2 从零搭起Spark集群:本地模式与Standalone的安装与使用

毕设阶段不推荐一上来就搭三台机器的YARN集群。大数据集群部署策略里有一条反复被验证的经验:先在单机把任务跑通,再谈分布式。本地模式适合写代码调试,Standalone适合演示"集群"效果,YARN留着等真有多个节点和HDFS需求时再上。下面的流程是spark的安装与使用里最常用的一套。

# 选Spark 3.x的稳定版预编译包,配置环境变量 export SPARK_HOME=/opt/spark export PATH=$SPARK_HOME/bin:$PATH # 单机Standalone:把默认配置写进spark-defaults.conf cat >> $SPARK_HOME/conf/spark-defaults.conf <<'EOF' spark.master spark://localhost:7077 spark.executor.memory 2g spark.driver.memory 1g spark.sql.shuffle.partitions 8 EOF # 启动Master和Worker $SPARK_HOME/sbin/start-master.sh $SPARK_HOME/sbin/start-worker.sh spark://localhost:7077 # 验证:进入spark-shell后能看到已连接的Worker $SPARK_HOME/bin/spark-shell --master spark://localhost:7077

逐个说参数含义。spark.master指向Standalone的Master地址,写成spark://格式可以在spark-submit和spark-shell里统一使用;spark.executor.memory是Executor的堆内内存,后面groupBy和join的数据都吃这块,毕设机器内存16G的话给2g比较稳;spark.sql.shuffle.partitions控制DataFrame触发shuffle时的分区数,默认200,单机小数据量会空转大量task,改成8能明显加快。

注意:spark.executor.memory不要超过机器可用内存的三分之二,否则Worker会频繁Full GC,表现是任务卡在某个stage不结束,日志里全是GC开销超限。

2.3 客流数据的三种接入方式与选型判断

数据接入层决定后面所有指标的口径。毕设里常见三种做法:直接读CSV、JDBC读MySQL、Kafka接实时流。直接读CSV最省事,AFC导出的刷卡流水本身就是CSV,不用搭额外组件;JDBC适合已经先做了一个管理前端、数据落在数据库里的情况;Kafka适合演示"大屏每分钟刷新"的效果,但需要同时维护生产者脚本和Streaming任务,调试成本最高。

接入方式数据形态演示效果实施成本
直接读CSV文件日结指标完整最低
JDBC读MySQL关系表与Web端共用数据源
Kafka + Structured Streaming消息流大屏分钟级刷新

我一般会先用CSV把整条链路打通,最后一周有余力再补Kafka。下面的CSV接入写法里,schema必须要提前声明,否则Spark会反复推断列类型,每次跑任务多花几十秒。

from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("metro-afc-etl") \ .master("local[4]") \ .config("spark.sql.shuffle.partitions", "8") \ .getOrCreate() afc_schema = """ card_id STRING, line_id STRING, station_id STRING, in_time TIMESTAMP, out_time TIMESTAMP, in_station STRING, out_station STRING, fare DOUBLE """ df = spark.read \ .option("header", "true") \ .option("timestampFormat", "yyyy-MM-dd HH:mm:ss") \ .schema(afc_schema) \ .csv("data/afc_2024_06_01.csv")

master里写local[4]表示本机4个线程并行,数据量小的时候比local[1]快4倍。timestampFormat必须和AFC导出的时间字符串完全一致,解析不了的字段会变null,这一步错了后面所有按时间的聚合都会缺数。

3. 用Spark DataFrame完成客流ETL与OD、断面指标计算

3.1 AFC原始数据的清洗规则与去重顺序

地铁客流数据看着整齐,实际脏数据集中在四个位置:闸机重复刷卡、进站无出站、时间倒挂、站名字段混入全角空格。清洗顺序会影响指标结果,所以先去重、再过滤、最后做字段归一,顺序不要反过来。先去重能降低后续groupBy和join的数据量;先过滤再去重,可能把本应合并的重复记录误删。

from pyspark.sql.functions import col, trim, regexp_replace # 1) 去重:同一张卡同一进站时间视为重复刷卡 dedup = df.dropDuplicates(["card_id", "in_time", "out_time"]) # 2) 过滤:无出站记录或时间倒挂的数据不参与OD计算 valid = dedup.filter( col("out_time").isNotNull() & (col("out_time") > col("in_time")) ) # 3) 归一:全角空格和普通空格统一去掉,避免同名站被当成两个站 clean = valid.withColumn( "in_station", trim(regexp_replace(col("in_station"), " |\\s+", "")) ).withColumn( "out_station", trim(regexp_replace(col("out_station"), " |\\s+", "")) )

dropDuplicates的subset要选够粒度,只按card_id去重会把一个人一天内的多段行程全部干掉,这是最常见的误用。时间倒挂的记录可以单独统计成"异常刷卡"口径,与正常客流分开计数,答辩时能讲清这个口径就是加分点。

3.2 站点进出量、OD矩阵、断面客流三个指标的DataFrame写法

三个核心指标的计算逻辑如下:站点进出量按"站+时间窗"聚合,OD矩阵按"进站+出站+时间窗"聚合,断面客流要把OD记录按线路展开成途经区间再聚合。前两个指标用groupBy就能完成,断面客流需要一张"线路站点顺序表"配合UDF展开。

from pyspark.sql.functions import window, count, col, explode, udf from pyspark.sql.types import ArrayType, StringType # 指标1:每半小时各站进站量 station_hour = clean.groupBy( window(col("in_time"), "30 minutes"), col("station_id") ).agg(count("card_id").alias("in_count")) # 指标2:OD矩阵,按小时聚合,用于桑基图和流向分析 od_matrix = clean.groupBy( col("in_station"), col("out_station"), window(col("in_time"), "60 minutes") ).agg(count("card_id").alias("od_flow")) # 指标3:断面客流,用线路站点表把OD展开成途经站 line_stations = { "line_1": ["station_001", "station_002", "station_003", "station_004"], "line_2": ["station_005", "station_006", "station_007"], } def stations_between(in_st, out_st): try: for seq in line_stations.values(): if in_st in seq and out_st in seq: i, j = seq.index(in_st), seq.index(out_st) return seq[i:j] if i < j else seq[j:i] except (KeyError, ValueError): return None return None expand_udf = udf(stations_between, ArrayType(StringType())) section_flow = clean \ .withColumn("pass_stations", expand_udf(col("in_station"), col("out_station"))) \ .select("line_id", "in_time", explode(col("pass_stations")).alias("pass_station")) \ .groupBy("line_id", "pass_station", window(col("in_time"), "30 minutes")) \ .agg(count("card_id").alias("section_flow"))

window函数里30 minutes是窗口大小,早高峰分析用30分钟粒度合适,换5 minutes数据量大5倍但能看出冲击细节。stations_between返回站名序列,explode展开成一条记录一个途经站,再做groupBy就得到断面流量。结果里station_001station_002区间的客流量,本质是所有经过该区间的OD记录数之和。UDF速度不算最快,但毕设数据量下完全够用,没必要为了性能引入复杂的高阶函数。

3.3 防止OOM和倾斜:spark内存与shuffle参数怎么设

客流数据里有一个天然倾斜点:早高峰的换乘大站。枢纽站进站量可能是普通站的三十倍,groupBy时这一个key会把某个task拖到超时。Spark 3.x的AQE能自动处理一部分倾斜,但需要先把开关打开。除此之外,下面三个参数是DataFrame任务里最常调的。

参数默认值建议值作用
spark.sql.shuffle.partitions200核心数×2~4控制shuffle输出分区数,过大会产生大量空task
spark.default.parallelism自动核心数×2~3影响stage初始并行度
spark.sql.adaptive.enabledfalsetrue开启动态合并与倾斜join优化
$SPARK_HOME/bin/spark-submit \ --master spark://localhost:7077 \ --executor-memory 2g \ --driver-memory 1g \ --conf spark.sql.shuffle.partitions=8 \ --conf spark.default.parallelism=8 \ --conf spark.sql.adaptive.enabled=true \ metro_etl.py

提交任务时我把参数全部写到spark-submit命令行而不是代码里,换机器跑批不用改源码。任务跑挂时先分两类日志:如果gc开销超过限制,说明Executor内存不够,往上加spark.executor.memory;如果某个task反复失败且伴随FetchFailedException,多半是数据倾斜,先确认AQE已开启,再检查执行计划里倾斜key的分区情况。spark内存相关的高频面试题集中在Executor堆内和堆外的划分,答辩时能说出"storage和execution共用一块内存池,互相可以抢占"这句,基本就说明理解到位了。

4. 客流分析结果落库到可视化大屏的完整链路

4.1 Spark结果写回MySQL:批量落库的参数控制

DataFrame计算完不能直接给前端用,需要落库。常见做法是Spark算完后写MySQL,Web后端只查MySQL,这样演示时前端不会因为Spark任务重跑而白屏。写JDBC时必须控制batchsize,一次性写入几十万行结果如果逐条插入,现场演示会在这一步卡住。

# 覆盖写OD结果表,batchsize控制单批次行数 od_matrix.write \ .mode("overwrite") \ .option("batchsize", "500") \ .option("truncate", "true") \ .jdbc( url="jdbc:mysql://localhost:3306/metro" "?useSSL=false&rewriteBatchedStatements=true", table="od_flow_result", properties={"user": "root", "password": "123456"} )

rewriteBatchedStatements=true是MySQL批量写入生效的关键,不加这个参数batchsize不生效,仍然逐条execute。mode选overwrite适合日结全量场景;如果做增量,改成append并配合日期分区字段,防止同一天的数据写重。结果表的字段要和DataFrame列名一致,否则JDBC映射报错,提前用printSchema()核对一遍最省时间。

4.2 用Spring Boot提供客流查询接口

后端接口尽量薄,只做"查表→返回JSON",不要在Controller里写聚合逻辑。聚合已经在Spark层做完了,后端重复算一是慢,二是容易和Spark口径不一致。下面的接口按站点和日期查分时进站量,用JdbcTemplate直接查结果表。

@RestController @RequestMapping("/api/flow") public class FlowController { private final JdbcTemplate jdbcTemplate; public FlowController(JdbcTemplate jdbcTemplate) { this.jdbcTemplate = jdbcTemplate; } @GetMapping("/station/hour") public List<Map<String, Object>> stationHour( @RequestParam String stationId, @RequestParam String date) { String sql = """ SELECT time_bucket, in_count FROM station_hour_flow WHERE station_id = ? AND stat_date = ? ORDER BY time_bucket """; return jdbcTemplate.queryForList(sql, stationId, date); } }

参数用?占位符绑定,避免拼接SQL引入注入风险,答辩抽查时这里经常被问。station_idstat_date要建联合索引,否则大屏轮询时接口会随结果表变大而越来越慢。如果前端刷新频率到了每秒一次,给接口加Spring Cache并配置5分钟过期,能挡住大部分重复查询;还要避免在循环里逐条查库,那是典型的N+1查询,接口会被越刷越慢。

4.3 数据大屏的三类图表与ECharts对接

大屏不要贪多,做三个能自圆其说的模块就够:全网总览、分时趋势、OD流向。热力站点图在地图API可用时是加分项,但地图服务的密钥配置和站点坐标数据准备会占大量时间,不作为必做项。下面是大屏模块与后端接口的对应关系。

大屏模块图表类型对应接口刷新方式
全网进出站总览数字翻牌+柱状图/api/flow/summary小时级
分时进站趋势面积折线图/api/flow/station/hour分钟级
OD客流量Top10桑基图/api/flow/od/top小时级

前端用ECharts时,关键是把后端返回的数组结构对齐到xAxis和series。下面这段是分时趋势图的加载逻辑,轮询间隔和后端缓存过期时间对齐,避免前端刷太快打满MySQL连接。

async function loadStationTrend(stationId, date) { const resp = await fetch( `/api/flow/station/hour?stationId=${stationId}&date=${date}` ); const rows = await resp.json(); chart.setOption({ xAxis: { type: 'category', data: rows.map(r => r.time_bucket) }, yAxis: { type: 'value', name: '进站量' }, series: [{ type: 'line', smooth: true, areaStyle: {}, data: rows.map(r => r.in_count) }] }); } // 大屏打开后每5分钟拉一次,与后端缓存过期时间对齐 loadStationTrend('station_101', '2024-06-01'); setInterval( () => loadStationTrend('station_101', '2024-06-01'), 5 * 60 * 1000 );

如果演示时想做出"实时感",可以加一个只刷新最近5分钟进站量的轻接口,而不是整屏重绘。把time_bucket排序后只取最后一条渲染增量,视觉上就是大屏在跳动的实时客流。用React或Vue写大屏时接口结构不变,只是把fetch逻辑封装进hooks或组合式函数里,后端完全不用动。

5. 答辩前必做的数据对账与Spark调优验证

5.1 用固定种子生成可复现的模拟客流数据

真实AFC数据涉及隐私,演示时换模拟数据是常见做法。生成模拟数据最关键的细节是固定随机种子,这样每次运行生成的CSV完全一致,Spark结果、MySQL表、大屏截图三者能对上号。

import random import pandas as pd random.seed(42) # 固定种子保证演示可复现 stations = [f"station_{i:03d}" for i in range(1, 21)] rows = [] for _ in range(20000): in_st = random.choice(stations) out_st = random.choice([s for s in stations if s != in_st]) h = random.randint(6, 23) rows.append({ "card_id": f"C{random.randint(10000, 99999)}", "in_time": f"2024-06-01 {h:02d}:{random.randint(0, 59):02d}:00", "out_time": f"2024-06-01 {h + 1:02d}:{random.randint(0, 59):02d}:00", "in_station": in_st, "out_station": out_st, "fare": random.choice([3.0, 4.0, 5.0]) }) pd.DataFrame(rows).to_csv("demo_afc.csv", index=False)

生成的字段与2.3节的schema对齐,放进data目录就能直接跑通ETL。20万条数据本地生成只需几秒,重跑代价低,适合边调参边验证。

5.2 用MySQL交叉验证Spark聚合结果

答辩时"结果对不对"比"怎么算"更难答。交叉验证做法是抽一个时间窗加一两个站点,把Spark输出表与对原始数据直接做的SQL统计对比。

SELECT station_id, '08:00-08:30' AS time_bucket, COUNT(*) AS in_count FROM afc_raw WHERE station_id = 'station_003' AND in_time >= '2024-06-01 08:00:00' AND in_time < '2024-06-01 08:30:00' GROUP BY station_id;

对比时两个口径必须对齐:Spark侧先做了dropDuplicates和空值过滤,MySQL侧也要先做同样的去重过滤,否则两边数字必然对不上。对账通过后把SQL结果和Spark输出记录存成截图,答辩现场直接展示,比口头解释"应该是对的"有说服力得多。

5.3 用checkpoint和结果表缓存兜底演示现场

演示最容易翻车的是Spark任务中途失败要重跑。checkpoint可以把中间stage结果落到磁盘,长链路跑到第三个stage挂掉时不用从头再算。调用方式很简单,但位置有讲究:

spark.sparkContext.setCheckpointDir("file:///tmp/spark_checkpoint") od_matrix = od_matrix.checkpoint() section_flow = section_flow.checkpoint()

注意checkpoint会切断血缘,不要每加一步都调用;只在链路中最长、最贵的那次shuffle之后checkpoint一次,能显著缩短失败恢复时间。大屏侧依赖的是MySQL结果表而不是Spark里的DataFrame,所以Spark重跑不影响前端展示,这才是整套系统演示时最稳的兜底设计。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询