ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

基于Spark的交通智能分析系统:从卡口流水到拥堵指数与OD矩阵

基于Spark的交通智能分析系统:从卡口流水到拥堵指数与OD矩阵 简介这份资源是面向高校学生与大数据初学者的一套完整项目源码以Apache Spark为核心构建交通智能分析系统适合用作毕业设计、课程作业或Spark实战练手。项目覆盖数据采集、预处理、交通流量分析、异常检测与决策支持等模块并延伸出电商用户行为分析的跨领域思路。压缩包共339个文件约1.45MB包含163个dat数据文件、129个class编译文件、13个scala与8个java源码以及xml配置、properties参数文件和少量txt说明完整保留了从源码到运行产物的工程结构。目前已有111人学习下载。通过研读这些代码读者可以掌握Spark Streaming实时接收GPS与摄像头数据、DataFrame清洗转换、Spark SQL聚合统计、MLlib流量预测等关键实现理解流式告警与状态监控的落地方式并借鉴其模块划分与目录组织快速搭建自己的大数据分析项目。1. 从一份交通卡口流水说起Spark 交通智能分析系统到底在算什么早高峰的卡口流水一天能到千万行字段无非是过车时间、设备编号、车牌哈希、车道号、车型。真正让人头疼的不是数据量而是同一辆车在 5 分钟内被 3 个相邻卡口拍到到底算一次出行还是三次。基于 Spark 的交通智能分析系统核心就是把这堆流水变成能回答问题的指标路段平均速度、拥堵指数、OD 出行矩阵、异常停留车辆。它适合两类人一类是手里已经有卡口或 GPS 数据、想搭一套离线准实时分析管道的工程师另一类是做课程设计或大数据实战需要一套能跑通、能讲清楚链路的完整项目。Spark 在这里的价值不是快而是把清洗、聚合、时空关联这几步用同一套 API 串起来批流代码几乎不用改。下面按我实际搭过的顺序从集群到指标一路拆开。2. 集群与数据落地Spark 交通分析系统的环境底座怎么搭2.1 为什么交通场景优先选 Spark 而不是纯 SQL 数仓交通数据的典型特征是宽表不宽、行数极多、时间窗口密集。一条卡口记录十几个字段但一天几千万行且大量计算是按 5 分钟窗口、按路段分组的滚动聚合。纯 SQL 数仓在固定报表上没问题但一旦要做 OD 推导、轨迹拼接这类需要多次 shuffle 和自定义状态的计算SQL 表达起来就很别扭。Spark 的优势在于DataFrame API 能覆盖 90% 的聚合剩下 10% 用 RDD 或 UDF 兜底同一份业务逻辑foreachBatch里改几行就能从离线切到 Structured Streaming。我一般会这样分工——离线层用 Spark SQL 跑 T1 的全量指标准实时层用 Structured Streaming 消费 Kafka 做 5 分钟粒度的拥堵预警。两套代码共享同一份清洗函数维护成本低很多。选型上还有一点常被忽略交通数据里车牌、设备编号这类维度表会变设备迁移、车牌脱敏规则调整Spark 的广播变量 维表关联比在数仓里做拉链表更灵活改一次广播逻辑就能生效。2.2 本地与集群的最小可跑环境新手最容易卡在环境上。我的建议是先用本地模式把逻辑跑通再上集群。本地用 Spark 3.x JDK 8/11Python 侧装 pyspark 即可不需要 Hadoop。# 本地伪分布式直接用 pip 装的 pyspark不依赖 Hadoop pip install pyspark3.5.0 # 验证 python -c from pyspark.sql import SparkSession; sSparkSession.builder.master(local[*]).getOrCreate(); print(s.version)集群侧常见做法是 Standalone 或 YARN。Standalone 适合课程设计和中小规模YARN 适合已有 Hadoop 集群的生产环境。关键参数就三个spark.executor.memory、spark.executor.cores、spark.sql.shuffle.partitions。前两个决定单 executor 能吃多少数据第三个决定 shuffle 后并行度。提示spark.sql.shuffle.partitions默认 200交通数据按天聚合时经常出现200 个分区里 190 个只有几 KB的小文件问题按数据量调到 分区数 ≈ 单日数据量(GB) × 4 比较稳。2.3 数据分层从原始流水到可分析宽表我一般分三层落地全部用 Parquet 存储层级内容分区键用途ODS原始卡口流水字段原样保留dt回溯、重跑DWD清洗后过车记录补路段 ID、去重dt, road_id明细查询DWS5 分钟粒度路段聚合指标dt, road_id报表、预警分区键选dt而不是时间戳是因为交通分析绝大多数查询都带日期范围按天分区能让 Spark 直接做分区裁剪扫描量降一个数量级。DWD 层额外按road_id二级分区是为了路段级指标计算时避免全表扫描。3. 数据清洗与时空关联把卡口流水变成可信过车记录3.1 清洗的四个必做动作原始流水里脏数据比例不低常见的有设备时钟漂移导致时间戳乱序、重复上报、车牌字段为空、车型编码不在字典里。清洗顺序不能乱我固定按去重 → 时间规整 → 字段补全 → 异常过滤走。from pyspark.sql import functions as F from pyspark.sql.window import Window def clean_pass_records(df): # 1. 去重同一设备、同一车牌、同一秒视为重复上报 df df.dropDuplicates([device_id, plate_hash, pass_time]) # 2. 时间规整设备时钟漂移超过 5 分钟的记录打标不直接丢 df df.withColumn( time_drift, F.abs(F.unix_timestamp(pass_time) - F.unix_timestamp(report_time)) 300 ) # 3. 字段补全车型编码缺失时用 unknown 占位避免后续 join 丢行 df df.fillna({vehicle_type: unknown, lane_id: -1}) # 4. 异常过滤速度字段为负或超过 200km/h 的剔除 df df.filter((F.col(speed) 0) (F.col(speed) 200)) return df逻辑说明去重放在最前是因为重复记录会污染后续所有聚合时间漂移只打标不删除是因为漂移记录的时间虽不可信但车牌和设备信息仍可用于轨迹补全车型缺失用占位符而非丢弃是为了保证 OD 矩阵的行数稳定。参数上5 分钟漂移阈值和 200km/h 上限都是经验值前者对应设备 NTP 同步周期后者对应城市道路物理极限高速场景要放宽到 160 以上。3.2 用窗口函数做轨迹拼接轨迹拼接的本质是给同一辆车的过车记录按时间排序算出相邻两条之间的时空差。这一步用窗口函数最直观。def build_trajectory(df): w Window.partitionBy(plate_hash).orderBy(pass_time) df df.withColumn(prev_device, F.lag(device_id).over(w)) \ .withColumn(prev_time, F.lag(pass_time).over(w)) \ .withColumn(prev_road, F.lag(road_id).over(w)) # 相邻卡口时间差用于判断是否为同一次出行 df df.withColumn( gap_sec, F.unix_timestamp(pass_time) - F.unix_timestamp(prev_time) ) # 超过 30 分钟视为出行中断重新起一段 df df.withColumn( trip_id, F.sum(F.when(F.col(gap_sec) 1800, 1).otherwise(0)).over(w) ) return df逻辑说明lag取上一条记录gap_sec是相邻过车的时间差。trip_id用累加的方式生成——每当 gap 超过 30 分钟就加 1这样同一辆车的一天会被切成若干段独立出行。参数 1800 秒是出行中断阈值通勤场景常用 30 分钟物流车辆可能要调到 2 小时。这里有个坑plate_hash如果脱敏规则变了同一辆车会被拆成多个 hash轨迹直接断掉所以脱敏规则必须版本化。3.3 路段匹配从设备编号到路段 ID卡口设备是点分析要的是路段。常见做法是维护一张设备-路段映射表用广播变量关联。from pyspark.sql.functions import broadcast device_road spark.read.parquet(/dim/device_road) # device_id, road_id, direction dwd df.join(broadcast(device_road), ondevice_id, howleft) # 未匹配到路段的记录单独落盘用于排查设备台账缺失 unmatched dwd.filter(F.col(road_id).isNull()) unmatched.write.mode(overwrite).parquet(/tmp/unmatched_device)逻辑说明映射表通常只有几千行用broadcast避免 shuffle。howleft保证不丢记录未匹配的单独落盘——这是排查设备台账问题的关键我见过太多项目直接 inner join结果新装设备的数据全被静默丢掉。参数上映射表要带direction字段因为同一路段上下行速度差异很大不区分方向算出来的平均速度没有意义。4. 核心指标计算拥堵指数、OD 矩阵与异常停留4.1 5 分钟粒度路段速度与拥堵指数路段速度的计算逻辑是路段长度 / 平均通行时间但直接用单车的点速度平均会失真因为慢车和快车权重一样。更稳的做法是用调和平均。def road_speed_5min(dwd): df dwd.groupBy( F.window(pass_time, 5 minutes).alias(win), road_id, direction ).agg( F.count(*).alias(flow), # 调和平均等价于 总距离/总时间对慢车更敏感 F.expr(count(*) / sum(1.0 / speed)).alias(avg_speed), F.expr(percentile_approx(speed, 0.85)).alias(p85_speed) ) # 拥堵指数 自由流速度 / 实际速度自由流速度取凌晨 3-5 点的 p85 df df.withColumn(congestion_index, F.col(free_flow_speed) / F.col(avg_speed)) return df逻辑说明调和平均count/sum(1/speed)在物理上就是总里程除以总时间比算术平均更贴近真实通行体验。p85_speed是第 85 百分位速度用来剔除极端慢车。拥堵指数用自由流速度除以实际速度自由流速度一般取该路段凌晨 3-5 点的 p85 速度这个值要离线算好存成维表。参数上5 分钟窗口是交通行业的通用粒度太细噪声大太粗丢失拥堵突变。4.2 OD 出行矩阵的推导ODOrigin-Destination矩阵是交通规划的核心输入本质是统计每个起点到每个终点的出行量。用前面拼好的trip_id聚合即可。def build_od_matrix(traj): trips traj.groupBy(plate_hash, trip_id).agg( F.first(road_id).alias(origin), F.last(road_id).alias(destination), F.min(pass_time).alias(start_time), F.max(pass_time).alias(end_time) ) # 过滤掉原地打转和过短的出行 trips trips.filter( (F.col(origin) ! F.col(destination)) (F.unix_timestamp(end_time) - F.unix_timestamp(start_time) 120) ) od trips.groupBy(origin, destination).agg(F.count(*).alias(trip_count)) return od逻辑说明first/last配合trip_id拿到每段出行的起终点。过滤条件里起终点相同的是原地折返可能是设备误拍时长小于 120 秒的出行大概率是相邻卡口误匹配。参数上120 秒阈值对应最短合理出行时间市中心可以调到 60 秒郊区要放宽。OD 矩阵输出后一般会做稀疏化只保留 trip_count 超过阈值的组合否则矩阵太稀疏没法用。4.3 异常停留车辆识别异常停留是交通管理的高频需求比如识别在应急车道长时间停留的车辆。逻辑是同一车辆在同一路段停留超过阈值。def detect_abnormal_stop(traj, threshold_sec600): stops traj.filter(F.col(gap_sec) threshold_sec) \ .select(plate_hash, prev_road, prev_time, pass_time, gap_sec) # 同一路段多次停留合并 stops stops.withColumn(stop_minutes, F.col(gap_sec) / 60) return stops逻辑说明gap_sec超过阈值说明车辆在相邻两次过车之间停留了很久prev_road就是停留路段。参数 600 秒是应急车道停留的常用判定阈值实际项目里会结合路段类型动态调整——高速应急车道 5 分钟就该报警普通路段 30 分钟才算异常。这里要注意设备故障导致的假停留很常见所以输出结果一般会再和该设备的在线状态做一次关联过滤。5. 避坑与排查Spark 交通分析系统最容易翻车的五个地方5.1 数据倾斜某个路段的数据量是其他的 100 倍现象任务卡在某个 stage 的 199/200日志显示某个 task 处理的数据量远超其他。原因热门路段比如市中心主干道的过车量天然是郊区的几十上百倍groupBy(road_id)时全压到一个分区。解决先对倾斜 key 加随机前缀打散聚合一次后再去前缀二次聚合或者直接开 AQEspark.sql.adaptive.enabledtrue让 Spark 自动处理倾斜 join。我一般两个都上AQE 兜底手动打散处理极端情况。5.2 时间窗口错位跨天数据被切成两半现象凌晨 0 点前后的 5 分钟窗口指标异常流量忽高忽低。原因window函数按自然时间切分跨天时分区裁剪把数据分到了两个dt分区但窗口本身跨了天。解决清洗时把pass_time统一转成 UTC 或固定时区且窗口计算前先按dt过滤再开窗避免跨分区。更稳的做法是窗口起点对齐到 5 分钟的整数倍而不是从第一条记录开始。5.3 车牌脱敏规则变更导致轨迹断裂现象某天开始 OD 矩阵的出行量骤降轨迹拼接结果里大量单点出行。原因上游脱敏规则从 MD5 换成了 SHA256同一辆车的plate_hash变了窗口函数按新 hash 分组历史轨迹全断。解决脱敏规则必须版本化hash 字段带版本前缀规则变更时要么重刷历史数据要么在拼接时用设备时间做兜底关联。这个坑我在两个项目里都踩过血泪经验是——脱敏字段永远不要当主键用。5.4 executor 内存溢出广播变量太大现象任务报OutOfMemoryError堆栈指向 broadcast。原因设备-路段映射表看着小但如果把整个维表含历史版本都广播几百万行直接撑爆 executor。解决广播前先过滤到当前有效版本控制广播变量在几十 MB 以内超过阈值就改用普通 join别硬广播。参数上spark.sql.autoBroadcastJoinThreshold默认 10MB交通维表经常超要盯着。5.5 小文件问题一天跑出几万个 Parquet 文件现象DWS 层目录下每个分区几万个几十 KB 的小文件下游查询慢NameNode 压力大。原因spark.sql.shuffle.partitions没调或者按高基数 key如plate_hash分区写入。解决写入前用repartition(n, dt)控制文件数n 按单分区数据量算一般 128MB 一个文件或者跑完用OPTIMIZEDelta或定时 compaction 合并。我习惯在写入 DWS 前统一coalesce到合理分区数简单有效。6. 进阶技巧用 Structured Streaming 把离线指标搬到准实时离线跑通之后真正有价值的是把拥堵指数做到分钟级延迟。我的做法是复用离线清洗函数只改数据源和输出。# 离线清洗函数直接复用保证批流逻辑一致 def process_batch(df, epoch_id): cleaned clean_pass_records(df) traj build_trajectory(cleaned) metrics road_speed_5min(traj) metrics.write.mode(append).parquet(/dws/road_speed_realtime) stream spark.readStream.format(kafka) \ .option(kafka.bootstrap.servers, localhost:9092) \ .option(subscribe, pass_records) \ .load() query stream.writeStream \ .foreachBatch(process_batch) \ .option(checkpointLocation, /checkpoint/road_speed) \ .trigger(processingTime1 minute) \ .start()逻辑说明foreachBatch让流处理复用批处理的 DataFrame 逻辑这是批流一体最实用的地方。checkpointLocation必须设置否则重启后状态丢失。trigger用 1 分钟微批兼顾延迟和吞吐。参数上Kafka 的maxOffsetsPerTrigger要设防止积压时一次拉太多数据把内存打爆。验证流处理是否正确我的习惯是拿同一天的离线结果和实时结果做对账按road_id和 5 分钟窗口 join比对avg_speed的差异超过 5% 就说明批流逻辑有偏差。这个对账脚本我一般会固化成日常任务比看日志靠谱得多。最后说个习惯交通数据的时区和节假日是两个隐形杀手。所有时间字段统一存 UTC展示层再转本地时区节假日的工作日/非工作日标记要单独维护一张日历表别用dayofweek硬判断。这两条我踩过坑之后写进了团队规范希望你也能少走弯路。希望帮到你。本文还有配套的精品资源点击获取
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进