简介:这是一份面向高校学生与大数据初学者的Spark实战项目资料,以交通智能分析系统为背景,适合用作毕业设计、课程作业或大数据入门练手。项目围绕数据采集、预处理、流量分析、异常检测与决策支持五个模块展开,借助Spark Streaming、Spark SQL与MLlib完成实时车流统计、拥堵预警和流量预测,其分析思路也可迁移到电商用户行为挖掘等场景。压缩包共339个文件,约1.45MB,包含163个dat数据文件、129个class编译文件,以及13个scala与8个java源码,另有xml配置、properties参数文件和少量txt说明,覆盖从源码到运行数据的完整结构。目前已有111人学习下载。通过研读源码与数据文件,读者可掌握Spark分布式计算、实时流处理与机器学习建模的落地方法,理解交通领域大数据系统的设计脉络,并积累可复用的项目经验与排错思路。
1. 从一份交通卡口流水说起:Spark 智能分析系统到底在算什么
早高峰的卡口流水一天能到千万行,字段无非是过车时间、设备编号、车牌哈希、车道号、车型。单机 pandas 读到三百万行就开始喘,groupby 一跑内存直接爆掉,这是很多人第一次意识到需要 Spark 的时刻。基于 Spark 的交通智能分析系统,本质就是把这类流水做成可查询、可聚合、可预警的指标层:拥堵指数、路段平均车速、早晚高峰流量对比、异常过车识别。它适合两类人,一类是要交课程设计或工程实践、需要一套能跑起来的完整链路;另一类是已经有一批卡口或 GPS 数据、想用 Spark 把离线统计和准实时聚合搭出来。标题里的「设计与实现」不是写文档,而是把数据从原始 CSV 一路推到能出图、能出报表的结果表。下面按我实际搭过的顺序拆开讲,先讲清楚数据怎么进来、怎么分层,再落到 Spark SQL 和调优参数,最后把踩过的坑摆出来。
2. 数据分层与表结构:交通流水的 ODS、DWD、DWS 怎么切
2.1 为什么不能一张宽表跑到底
交通数据的天然特征是「原始流水极宽、指标极窄」。原始过车记录有几十个字段,但真正用于分析的往往只有时间、路段、方向、车型四五个维度。如果所有计算都直接怼在原始表上,每次跑指标都要全量扫一遍,磁盘 IO 和 shuffle 都吃不消。常见做法是分三层:ODS 层保留原始接入数据,字段和来源一致,只做格式统一;DWD 层做清洗和标准化,把时间戳归一到秒、把设备编号映射到路段、把无效车牌和重复过车剔掉;DWS 层按「路段 + 时间窗」预聚合,直接产出流量、平均车速、饱和度这类指标。这样上层做报表或预警时只扫 DWS,数据量能降一到两个数量级。
分层还有一个隐性好处:排错时能定位到具体环节。指标不对,先看 DWD 的清洗规则,再看 DWS 的聚合口径,不用在一张大宽表里猜是哪一步出的问题。我一般会把每层的分区字段统一成dt(日期)加hour(小时),这样按天回溯和按小时补数都方便。
2.2 建表语句与分区设计
下面这套 DDL 是我在 Hive 外表上常用的结构,ODS 用文本或 Parquet 都行,DWD 和 DWS 统一 Parquet 加 snappy 压缩。
-- ODS:原始过车流水,按天分区 CREATE TABLE ods_traffic_pass ( pass_time STRING COMMENT '过车时间 yyyy-MM-dd HH:mm:ss', device_id STRING COMMENT '卡口设备编号', plate_hash STRING COMMENT '车牌哈希,脱敏后', lane_no INT COMMENT '车道号', vehicle_type STRING COMMENT '车型', speed DOUBLE COMMENT '瞬时车速,可能为空' ) PARTITIONED BY (dt STRING) STORED AS PARQUET; -- DWD:清洗后标准流水,补上路段维度 CREATE TABLE dwd_traffic_pass ( pass_ts BIGINT COMMENT '过车时间戳,秒', road_id STRING COMMENT '路段编号', direction STRING COMMENT '方向:上行/下行', plate_hash STRING, lane_no INT, vehicle_type STRING, speed DOUBLE ) PARTITIONED BY (dt STRING, hour STRING) STORED AS PARQUET; -- DWS:路段小时级指标 CREATE TABLE dws_road_hour_metric ( road_id STRING, direction STRING, hour STRING, pass_cnt BIGINT COMMENT '过车数', avg_speed DOUBLE COMMENT '平均车速', congestion DOUBLE COMMENT '拥堵指数 0-10' ) PARTITIONED BY (dt STRING) STORED AS PARQUET;分区字段的选择直接决定后续查询能不能裁剪。dt放最外层是为了按天批量重跑,hour放内层是为了小时级补数。注意 DWD 的pass_ts用 BIGINT 而不是字符串,后面做时间窗聚合时省掉一次unix_timestamp转换,这个细节在千万级数据上能省不少 CPU。
2.3 设备到路段的映射表怎么维护
卡口设备编号和路段的对应关系不是一成不变的,新建设备、临时改道都会让映射失效。我一般单独建一张维表dim_device_road,字段是device_id、road_id、direction、start_date、end_date,用拉链方式保留历史。DWD 清洗时按pass_time落在start_date和end_date之间去 join,避免用最新映射去套历史数据导致路段归属错乱。这张表数据量小,可以广播到每个 executor,join 时不会有 shuffle。
3. 用 Spark SQL 把原始流水跑成拥堵指标
3.1 从 ODS 到 DWD 的清洗逻辑
清洗这一步的核心是去重、补维、算时间戳。去重不能简单按车牌去,同一辆车短时间内多次过同一设备可能是重复上报,也可能是真的绕了一圈,我一般按「设备 + 车牌 + 分钟」做窗口去重,保留最早一条。
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window spark = SparkSession.builder \ .appName("traffic_dwd_clean") \ .config("spark.sql.shuffle.partitions", "200") \ .config("spark.sql.adaptive.enabled", "true") \ .enableHiveSupport() \ .getOrCreate() # 读 ODS 当天分区 ods = spark.table("ods_traffic_pass").filter(F.col("dt") == "2024-06-01") # 时间戳归一 + 分钟窗口去重 ods = ods.withColumn("pass_ts", F.unix_timestamp("pass_time", "yyyy-MM-dd HH:mm:ss")) w = Window.partitionBy("device_id", "plate_hash", F.floor(F.col("pass_ts") / 60)).orderBy("pass_ts") dwd = ods.withColumn("rn", F.row_number().over(w)) \ .filter(F.col("rn") == 1) \ .drop("rn") # 广播维表补路段 dim = spark.table("dim_device_road") \ .filter(F.col("end_date").isNull() | (F.col("end_date") >= "2024-06-01")) dwd = dwd.join(F.broadcast(dim), on="device_id", how="left") \ .withColumn("hour", F.lpad(F.hour("pass_time"), 2, "0")) \ .select("pass_ts", "road_id", "direction", "plate_hash", "lane_no", "vehicle_type", "speed", "dt", "hour") dwd.write.mode("overwrite").insertInto("dwd_traffic_pass")spark.sql.shuffle.partitions设 200 是经验值,数据量在几千万行时比较稳,太小会 OOM,太大会产生大量小文件。spark.sql.adaptive.enabled打开后 Spark 会自动合并小分区,对倾斜的卡口数据特别有用——某些主干道设备的数据量可能是支路的几十倍,自适应能缓解长尾。广播维表用F.broadcast显式提示,避免优化器判断失误走 shuffle join。
3.2 小时级指标聚合与拥堵指数计算
DWS 层的聚合是整套系统的核心。过车数直接 count,平均车速用avg(speed),但要注意空值处理——很多卡口不返回瞬时车速,直接 avg 会把空值算进去导致结果偏低,我一般先过滤speed is not null再算,同时记录有效样本数。
拥堵指数没有统一公式,常见做法是用「实际行程时间 / 自由流行程时间」的比值再映射到 0-10。没有行程时间数据时,可以用平均车速反推:设自由流车速为v_free(比如 60 km/h),拥堵指数 =min(10, v_free / avg_speed)。这个公式粗糙但可解释,适合课程设计或初期版本。
dwd = spark.table("dwd_traffic_pass").filter(F.col("dt") == "2024-06-01") dws = dwd.groupBy("road_id", "direction", "hour", "dt").agg( F.count("*").alias("pass_cnt"), F.avg(F.when(F.col("speed").isNotNull(), F.col("speed"))).alias("avg_speed"), F.sum(F.when(F.col("speed").isNotNull(), 1).otherwise(0)).alias("speed_samples") ) # 拥堵指数:自由流 60,无有效车速时置空 dws = dws.withColumn( "congestion", F.when(F.col("avg_speed") > 0, F.least(F.lit(10.0), F.lit(60.0) / F.col("avg_speed"))) .otherwise(None) ) dws.write.mode("overwrite").insertInto("dws_road_hour_metric")F.least用来封顶,避免低速时指数飙到几十。speed_samples这个字段很多人会漏,它的作用是让下游知道这条指标可不可信——样本数只有个位数时,平均车速波动极大,报表里应该标注出来。
3.3 早晚高峰对比怎么查
有了 DWS,早晚高峰对比就是一句 SQL 的事。下面查每个路段早高峰(7-9 点)和晚高峰(17-19 点)的流量与平均车速对比。
SELECT road_id, direction, SUM(CASE WHEN hour BETWEEN '07' AND '09' THEN pass_cnt ELSE 0 END) AS am_cnt, SUM(CASE WHEN hour BETWEEN '17' AND '19' THEN pass_cnt ELSE 0 END) AS pm_cnt, AVG(CASE WHEN hour BETWEEN '07' AND '09' THEN avg_speed END) AS am_speed, AVG(CASE WHEN hour BETWEEN '17' AND '19' THEN avg_speed END) AS pm_speed FROM dws_road_hour_metric WHERE dt = '2024-06-01' GROUP BY road_id, direction ORDER BY am_cnt DESC LIMIT 50;这里用CASE WHEN而不是两次查询再 join,减少一次扫描。AVG对avg_speed再平均是近似值,严格来说应该用加权平均(按 pass_cnt 加权),如果对精度要求高,把SUM(avg_speed * pass_cnt) / SUM(pass_cnt)换上去即可。
4. 避坑与排查:交通数据在 Spark 上最容易翻车的五件事
4.1 数据倾斜:某个卡口的数据量是别人的一百倍
现象是任务卡在最后一个 reduce 阶段不动,打开 Spark UI 看到某个 task 处理的数据量远超其他。原因是主干道卡口的过车量天然远高于支路,groupBy 路段时数据全挤到少数分区。解决办法分两步:先确认倾斜键,用dwd.groupBy("road_id").count().orderBy(F.desc("count")).show(10)找出来;再对倾斜键加随机前缀打散,聚合完再去掉前缀二次聚合。如果用的是 Spark 3.x,直接开 AQE 的倾斜处理spark.sql.adaptive.skewJoin.enabled=true也能缓解大部分场景。
4.2 时间戳时区错乱导致高峰时段偏移
现象是早高峰统计算出来落在凌晨。原因是unix_timestamp默认按 JVM 时区解析,集群时区是 UTC 时,yyyy-MM-dd HH:mm:ss会被当成 UTC 时间,和本地时间差 8 小时。解决方式是在 SparkSession 里显式设spark.sql.session.timeZone=Asia/Shanghai,或者清洗时统一用to_timestamp并指定格式,不要依赖默认行为。这个坑很隐蔽,因为数据本身没错,错的是解释方式。
4.3 小文件过多拖垮 NameNode
现象是跑完一天数据后 HDFS 上多出几万个几十 KB 的文件,下次查询光列目录就要好几秒。原因是分区粒度太细(dt + hour + 设备)加上并行度高。解决办法是写入前用repartition控制文件数,或者在 DWS 层按 dt 聚合后只写一个分区。我一般会在写入后跑一个合并脚本,把小于 128MB 的文件合并掉。注意coalesce和repartition的区别:前者不 shuffle 但可能造成分区不均,后者 shuffle 但分布均匀,按数据量选。
4.4 车牌哈希后仍然能反推
现象是脱敏做了但还能通过哈希碰撞或字典攻击还原。原因是用了无盐的 MD5 或 SHA1,车牌空间小,彩虹表一查就出来。解决方式是加固定盐值再哈希,或者直接用 HMAC。盐值不要写在代码里,放配置中心或环境变量。这个坑在课程设计里经常被忽略,但一旦数据外流就是实打实的隐私问题。
4.5 内存参数照搬网上配置导致 executor 被 kill
现象是任务跑一半报Container killed by YARN for exceeding memory limits。原因是spark.executor.memory设得大但spark.executor.memoryOverhead没跟上,或者spark.memory.fraction默认 0.6 导致执行内存不够。我的习惯是 executor 内存设 4-8G,overhead 给 10%-15%,spark.memory.fraction保持默认,先跑小数据量压测再放大。另外spark.sql.shuffle.partitions和 executor 数量要匹配,200 个分区配 4 个 executor 就是每个 executor 扛 50 个分区,容易 OOM。
5. 把离线指标接到准实时:Structured Streaming 的增量聚合技巧
离线跑通之后,很多人会想能不能让拥堵指数分钟级更新。Structured Streaming 是 Spark 原生的流处理入口,和批处理共用一套 API,迁移成本低。我的做法是把 Kafka 里的过车消息按「路段 + 5 分钟窗口」做增量聚合,输出到 DWS 的实时表,离线表负责 T+1 校准,两张表在查询层 union。
stream = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "broker:9092") \ .option("subscribe", "traffic_pass") \ .load() parsed = stream.selectExpr("CAST(value AS STRING)") \ .select(F.from_json("value", "pass_ts LONG, road_id STRING, direction STRING, speed DOUBLE").alias("d")) \ .select("d.*") \ .withWatermark("pass_ts", "2 minutes") agg = parsed.groupBy( F.col("road_id"), F.col("direction"), F.window(F.col("pass_ts").cast("timestamp"), "5 minutes") ).agg( F.count("*").alias("pass_cnt"), F.avg("speed").alias("avg_speed") ) agg.writeStream \ .outputMode("update") \ .format("parquet") \ .option("path", "/warehouse/dws_road_realtime") \ .option("checkpointLocation", "/checkpoint/traffic_realtime") \ .trigger(processingTime="1 minute") \ .start()withWatermark设 2 分钟是容忍乱序的边界,交通数据从设备上报到进 Kafka 一般延迟在秒级,2 分钟足够覆盖大部分乱序。outputMode("update")只输出变化的窗口,比complete省资源。trigger设 1 分钟意味着微批间隔,延迟和吞吐的折中点。checkpoint 目录必须放在可靠存储上,否则重启后状态丢失会重复计算。
这里有个容易忽略的点:流处理的聚合结果和离线结果对不上是正常的,因为窗口边界和乱序处理策略不同。我的习惯是在报表层标注数据来源,实时值用于监控告警,离线值用于正式统计,两者不混用。另外流式写入 parquet 会产生大量小文件,生产环境一般换成 Delta 或 Hudi,课程设计阶段用 parquet 加定时合并也能接受。
调优上,流处理对 executor 数量不敏感但对内存敏感,因为状态要常驻。spark.sql.streaming.stateStore.providerClass默认的 HDFSBackedStateStore 在状态大时性能一般,数据量上来后可以换 RocksDB。这些参数不用一开始就调,先让链路跑通,压测时再逐个改。
最后说个我自己的习惯:每次改完聚合逻辑,先拿一天的历史数据用批处理跑一遍,确认指标口径对了,再切到流。批流用同一套 SQL,能省掉大量对账时间。这套系统值不值得做,取决于你手上有没有持续产生的交通数据——有数据,Spark 的分层和聚合能力能让你从「导出一张 Excel」变成「随时查任意路段任意时段」;没数据,先把公开的卡口数据集跑通链路,再考虑接真实源。希望帮到你。
本文还有配套的精品资源,点击获取