简介:《基于Spark+Hive的交通智能研判系统》是一套用于毕业设计和课程设计场景的大数据实践项目,基于Spark与Hive两大组件构建,面向交通流量实时计算、历史数据管理和智能研判分析等问题,适合正在学习或开发分布式数据处理系统的读者参考。压缩包内共58个文件,以42个Java源文件为主体,另有9个XML配置、2个properties文件以及少量辅助文件,整体体积仅953KB,便于快速下载和本地运行。目前已有147人浏览/学习该资源,具有一定的参考价值。项目在实现中展示了Spark对大规模交通数据的快速读取与内存计算能力,以及Hive对历史数据的结构化存储与SQL查询能力,可覆盖车流量统计、拥堵原因分析、交通时段对比等典型场景;包内还提供Maven工程结构和IDE配置,便于导入开发工具后直接运行或二次改造,是一份能帮助答辩演示与技能提升的完整实践素材。
1. 拿到「交通智能研判系统」源码后,我第一件事不是跑起来
一个名为“基于Spark+Hive的交通智能研判系统”的压缩包,解压后是 TrafficTeach-master,里面有 pom.xml、src/main、monitor_camera_info、monitor_flow_action 这些目录和文件。很多同学下载这类毕设源码的第一反应是直接 mvn package 然后到处点,结果要么 Hive 连不上,要么 Spark 作业提交后卡死。我第一次看到这个工程时反而先做了一件事:把每个目录和两张核心表的关系捋清楚。Spark 负责吃实时的过车流水,Hive 负责沉淀历史数据供离线研判,两者通过元数据服务打通。这篇笔记的目标读者很明确:准备拿 Spark+Hive 做毕业设计或课程设计的人,以及想看看一个教学型大数据工程长什么样、哪些地方能直接复用的从业者。
2. 拆包 TrafficTeach-master:先搞懂 Spark 和 Hive 在工程里各管哪一段
2.1 目录与文件逐个过一遍
用 IDEA 打开工程之前,我习惯先把压缩包解压后的顶层文件看一遍。这个工程的结构不算复杂,但每个文件都有明确用途,我逐个说。
根目录下的 TrafficTeach-master 是工程主目录,里面能看到这些关键内容:
- pom.xml:Maven 工程描述文件,定义了 Spark、Hive 相关依赖。原则上,这个工程能直接用
mvn clean package构建,但前提是你本地 Maven 仓库里已经拉过这些依赖,第一次构建通常要等一段时间。 - src/main:源代码目录,Spark 的实时计算逻辑和离线分析逻辑都写在这里。
- src/test:单元测试目录,教学型工程里这一块通常不完整,但保留目录结构是给后面扩展留位置。
- .idea/:IntelliJ IDEA 的工程配置目录。有它说明作者开发时用的是 IDEA,你直接用 IDEA 打开 TrafficTeach-master 目录就能识别成 Maven 工程,省去手动配 SDK 的麻烦。当然,用 Eclipse 也不是不行,但 IDEA 对这个工程更友好。
- TrafficTeach.iml:IDEA 的模块描述文件,和 .idea 配套。
- target/:编译输出目录。压缩包里带 target 说明作者本地已经成功构建过一次,如果你打开后 target 里的 class 文件还在,说明这套代码在作者的机器上是能编译过的,这是一个“工程本身没问题”的强信号。
- monitor_camera_info 和 monitor_flow_action:这两项是数据文件或数据表名。前者是监控摄像头点位信息,属于维度表,记录摄像头编号、所在路口、方向、所属区域这类静态信息;后者是过车行为流水,属于事实表,记录每一辆车经过某个摄像头的时间、车牌、车型等动态信息。
我特别注意 monitor_camera_info 和 monitor_flow_action 这两个名字,因为整个交通智能研判系统的业务逻辑都围绕“摄像头”和“过车行为”展开。Spark 实时处理的是 flow_action 这条流水,Hive 离线分析时要 join 上 camera_info 才能把摄像头 ID 翻译成具体的区域和路口。
提示:拿到任何 Spark 教学工程,第一件事不是改代码,而是先找到“维度表”和“事实表”。维度表描述“是什么”,事实表记录“发生了什么”,这两张表一旦定位清楚,整个项目的业务流程就出来了一半。
2.2 Spark 和 Hive 的分工边界
很多初学者会把 Spark 和 Hive 混为一谈,觉得它们都是处理数据的,为什么要同时用两个。这个工程恰好把两者的分工讲得很清楚,我拆包后最大的收获也在这里。
Spark 在这个系统里负责的是“实时或近实时”的链路。过车流水一条条进来,Spark 用滑动窗口每隔几分钟算一次各卡口的车流量、平均车速等指标。这类计算要求延迟低,能对源源不断的数据做增量处理,这正是 Spark Structured Streaming 的强项。而 Hive 负责的是“历史沉淀”。过车流水会持续写入 Hive 的表中,日复一日地积攒下来。当你想看某个区上个月的日均车流、对比节假日和工作日的拥堵差异时,就需要用 HQL 去汇总这批历史数据。
两者的衔接点是 Hive 的元数据服务。Spark 作业在写入结果时,借助 Hive metastore 找到目标表并写入;Hive 查询时读到的也是同一套元数据。所以,只要hive-site.xml里的元数据库地址一致,Spark 和 Hive 就能共用一套表结构。如果哪个环节配错了,就会出现“Spark 写完了 Hive 查不到”这种经典问题,我后面专门讲。
用一张表概括这个工程的分工:
| 组件 | 数据角色 | 典型动作 | 时效性 |
|---|---|---|---|
| Spark Structured Streaming | 流式处理引擎 | 窗口聚合、实时研判 | 分钟级 |
| Hive | 数据仓库工具 | 历史存储、离线汇总 | 小时级/天级 |
| Hive metastore | 元数据中枢 | 打通 Spark 与 Hive 的表定义 | 即时 |
2.3 pom.xml 的依赖设计
我习惯先看 pom.xml 里的依赖,它能直接告诉我这个工程是 Java 写的还是 Scala 写的,以及 Spark 用的是哪个大版本。对于这类教学工程,最常见的组合是 Spark 2.4.x 或 3.x 加上对应的 spark-sql、spark-hive 依赖。
典型的依赖坐标长这样:
<dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_2.12</artifactId> <version>2.4.8</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-hive_2.12</artifactId> <version>2.4.8</version> <scope>provided</scope> </dependency>scope 设成 provided 意味着运行环境中已经提供了 Spark 的 jar 包,打完包后不会把整套 Spark runtime 塞进去,防止和你提交作业的集群环境产生依赖冲突。这个细节很多课设报告里压根不提,但实际集群部署时非常关键。
如果你想用本地模式直接跑这个工程做调试,记得把 provided 改成 compile 或者注释掉 scope,否则 IDE 里运行时会报找不到 spark session 相关类。这算是我拆这类工程遇到的第一个小坑,放在这里给后来人提个醒。
3. Spark 实时流计算:从“过车流水”到“卡口车流”的实战写法
3.1 实时链路的数据入口
教学工程通常不会真去接卡口摄像头的 Kafka 消息,最常见做法是用一个 JSON 或 CSV 目录模拟数据源,Spark 用readStream去读这个目录,然后按窗口聚合。实际生产里,这个位置多半换成 Kafka,但从 DataFrame 的写法来说,改动只在数据源格式那一行。
用文件目录模拟数据源时,我的标准写法是这样的:
from pyspark.sql import SparkSession from pyspark.sql.functions import window, col, count spark = SparkSession.builder \ .appName("TrafficFlowAnalysis") \ .master("local[2]") \ .enableHiveSupport() \ .config("spark.sql.warehouse.dir", "hdfs://namenode:9000/user/hive/warehouse") \ .getOrCreate() flow_stream = spark \ .readStream \ .format("json") \ .option("inferSchema", "true") \ .option("maxFilesPerTrigger", "1") \ .load("file:///data/traffic_flow")这段代码里有两个参数值得注意。maxFilesPerTrigger设成 1,意思是每次触发只读一个新文件,这样在本地模拟流式效果时比较平滑,不会一次性把所有数据全部灌进来,方便观察窗口计算的变化过程。enableHiveSupport()则是在 SparkSession 里打开 Hive 支持开关,没有这一行,后面想用 Spark 直接读写 Hive 表会直接报错。
如果工程本身是 Java 写的,只是语言从 Python 换成 Java API,流程完全一样,DataFrame 的算子名称也无差别,只是包名从pyspark.sql.functions变成org.apache.spark.sql.functions。这也是我建议这类毕设先跑通 Python 版本再回头看 Java 源码的原因:逻辑先成立,语言只是载体。
3.2 滑动窗口的参数到底怎么设
交通流研判最核心的指标之一,是“当前时间段内各卡口通过多少辆车”。这里要用到 Spark 的滑动窗口,而不是简单的 Tumbling Window。滑动窗口能避免每分钟的边界硬切,让统计结果更平滑,也更接近交通管理者的直觉。
窗口聚合的完整写法:
car_flow = flow_stream \ .withWatermark("pass_time", "5 minutes") \ .groupBy( window(col("pass_time"), "10 minutes", "5 minutes"), col("camera_id") ) \ .agg(count("vehicle_plate").alias("flow_count"))我拆这个工程时花了最多时间研究的就是这里。
window(col("pass_time"), "10 minutes", "5 minutes")的语义是:窗口长度为 10 分钟,每 5 分钟滑动一次。也就是说,整个时间轴会被切成长度为 10 分钟、彼此重叠 5 分钟的窗口序列。一辆车在 08:00:30 经过卡口,它会同时落入 [08:00, 08:10) 和 [08:05, 08:15) 这两个窗口,所以在原始计数逻辑下,同一辆车会被算两次。这不算 Bug,这是滑动窗口的固有特性。如果你要的是“不重复的去重车流量”,就得改用approx_count_distinct或者按车辆 ID 再做一次去重,而不是质疑窗口写错了。
withWatermark("pass_time", "5 minutes")解决的是乱序数据的问题。卡口设备上传过车记录时经常因为网络抖动延迟到达,比如 08:03 产生的记录 08:09 才进入流。如果这个延迟时间超过了 watermark 的容忍范围,这条数据就会被判定为“太晚到达”而丢弃。把 watermark 设成 5 分钟,意味着允许最多 5 分钟的迟到窗口,但代价是结果不会立刻输出,后端聚合会有最多 5 分钟的事件时间延迟。这块是 Spark Streaming 里最容易“凭感觉调参”的地方,我的建议是:先用自己的历史数据回放一遍,统计 90% 的记录延迟分布,再决定 watermark 取值,而不是拍脑袋写 5 分钟。
3.3 聚合结果写回 Hive
实时链路算出来的卡口车流是明细级的汇总值,要让它对业务产生价值,最终得落到 Hive 表里,让后续离线分析和报表查询能读到。Structured Streaming 写 Hive 有一个常用做法,叫foreachBatch,每次微批处理拿到的都是一个静态 DataFrame,可以直接用标准 DataFrame writer 写入 Hive。
写法示例如下:
def write_to_hive(batch_df, batch_id): batch_df.write \ .mode("append") \ .insertInto("dwd_traffic.traffic_flow_window") car_flow \ .writeStream \ .foreachBatch(write_to_hive) \ .outputMode("append") \ .trigger(processingTime="1 minute") \ .start() \ .awaitTermination()insertInto要求表已经存在于 Hive 中,并且批处理 DataFrame 的列顺序、列名要与目标表一致。这个坑特别隐蔽:字段名对不上时 Spark 并不报错,而是按位置匹配,结果就是数据写进去了,但你查出来发现 camera_id 对应的值其实是时间字段。所以我在写完这个逻辑后,一定会强制在 HQL 里跑一条SELECT * FROM dwd_traffic.traffic_flow_window LIMIT 5,肉眼核对一下列对应关系再继续。
trigger(processingTime="1 minute")是控制微批触发频率的,本地调试时可以放宽到 1 分钟,但如果是生产环境且数据量不大,5 分钟甚至 10 分钟一次会更省资源。这里的取舍在于“你对实时性的要求”和“你对集群成本的承受力”之间的平衡。
提示:foreachBatch 里的 batch_df 不要直接用
collect()拉回 Driver,否则数据量一上来,Driver 内存直接被撑爆。让数据在 Executor 端写完,这是写 Streaming 作业的基本习惯。
4. Hive 离线分析:建表、分区与交通态势对比查询
4.1 表结构设计:维度表、事实表与分区策略
Hive 侧的表设计直接决定离线查询的效率和可维护性。这个工程里的两张核心表,我建议按下面的方式建模。
摄像头点位维度表,用外部表存储,因为这类数据来自业务系统,Hive 只负责读:
CREATE EXTERNAL TABLE dim_traffic.monitor_camera_info ( camera_id STRING COMMENT '摄像头编号', road_name STRING COMMENT '所在道路', area_name STRING COMMENT '所属区域', direction STRING COMMENT '朝向:南北/东西', latitude DOUBLE, longitude DOUBLE ) ROW FORMAT DELIMITED FIELDS TERMINATED BY ',' STORED AS TEXTFILE LOCATION '/warehouse/traffic/dim/monitor_camera_info';过车行为事实表,按日期分区存储,这是离线分析里最重要的设计决策:
CREATE EXTERNAL TABLE dwd_traffic.traffic_flow_action ( camera_id STRING COMMENT '摄像头编号', vehicle_plate STRING COMMENT '车牌号', vehicle_type STRING COMMENT '车型:小车/货车/客车', pass_time STRING COMMENT '过车时间,格式 yyyy-MM-dd HH:mm:ss' ) PARTITIONED BY (dt STRING COMMENT '数据分区,格式 yyyy-MM-dd') STORED AS ORC LOCATION '/warehouse/traffic/dwd/traffic_flow_action';把 dt 设成分区字段,核心目的是查询裁剪。比如你要查 2024 年 6 月的交通态势,Hive 只需要读取dt >= '2024-06-01'对应的分区目录,而不必扫描全表。如果你用 Spark 写入这个表,记得在写之前按 dt 做一次分区写出,让每个 Spark 任务只负责自己那天的数据,避免产生大量跨区小文件。
关于数据格式,没用 TEXTFILE 而用 ORC,是因为 ORC 是列式存储。离线分析里最常见的查询是“按时间和区域维度聚合”,只涉及到其中少数几列,列式存储可以跳过无关列,查询速度差距非常明显。这也解释了为什么这个工程需要 Hive,而不是直接在 HDFS 上查原始文件。
4.2 工作日与周末的交通态势对比
摘要里特别提到“节假日与工作日的对比分析”,这正好是 Hive 离线分析的典型场景。核心思路是从事实表里把每天的过车记录按区域聚合,再分成工作日和周末两组做对比。
我常用的 HQL 如下:
SELECT c.area_name, IF(pmod(datediff(dt, '1970-01-01'), 7) < 5, 'workday', 'weekend') AS day_type, COUNT(*) AS total_flow, COUNT(DISTINCT a.vehicle_plate) AS distinct_vehicle_cnt, ROUND(AVG(a.flow_per_camera), 2) AS avg_flow_per_camera FROM ( SELECT camera_id, vehicle_plate, dt, COUNT(*) AS flow_per_camera FROM dwd_traffic.traffic_flow_action WHERE dt >= '2024-06-01' AND dt <= '2024-06-30' GROUP BY camera_id, vehicle_plate, dt ) a JOIN dim_traffic.monitor_camera_info c ON a.camera_id = c.camera_id GROUP BY c.area_name, IF(pmod(datediff(dt, '1970-01-01'), 7) < 5, 'workday', 'weekend');这里的 IF 判断是计算一个日期是工作日还是周末的关键。datediff(dt, '1970-01-01')得到从纪元到当前的天数,再pmod(... , 7)取模,结果是 0 到 6,分别对应周日到周六。小于 5 即为周一到周五,else 分支就是周六周日。这个写法不依赖任何自定义函数,纯 HQL 就能完成,比在 Spark 里先算好再导过来省事得多。
嵌套子查询的作用是先去一行一行地聚合,得到“每辆车在每个摄像头每天通过几次”,再在这个粒度上做区域汇总。这样一个区域一天的总车流,等于所有摄像头当天的总过车次数;去重车数则是 DISTINCT 车牌数。AVG(flow_per_camera) 才是真正的“单摄像头平均过车量”,比单纯看总量更能反映拥堵程度。
如果你只想看某个具体路口的早晚高峰差异,把 GROUP BY 粒度改为c.road_name再加上HOUR(pass_time)分组即可。这也是我在做课设时经常被问到的改法:把维度从区域换成路口,把时间从整天换成小时,查询结构不变,改两个字段就行。
4.3 分区裁剪与 Hive 查询的隐性性能门槛
新手最容易犯的错误是:分区表建好了,但 HQL 里不写分区过滤条件,或者写成函数嵌套形式导致分区裁剪失效。
举个例子:
-- 错误写法,分区裁剪失效 SELECT * FROM dwd_traffic.traffic_flow_action WHERE dt >= SUBSTR('2024-06-30', 1, 10); -- 正确写法,分区裁剪生效 SELECT * FROM dwd_traffic.traffic_flow_action WHERE dt >= '2024-06-30';第一种写法里,SUBSTR('2024-06-30', 1, 10)的结果虽然是 '2024-06-30',但 Hive 的优化器无法在执行前确定它的值是否恒定,干脆放弃分区裁剪,直接全表扫描。对于一些数据量只有几万条的教学数据集,两种写法看不出差别;但如果你想拿这个工程参加答辩,把数据量扩到几百万条再演示,有分区和没分区的查询时间可能是秒级与分钟级的区别。
另一个经常被忽略的问题是“小文件过多”。如果 Spark 写入时用了过细的分区粒度,比如每 5 分钟一个分区,一天的 Hive 表可能就有几百个小文件,NameNode 内存和查询效率都会受影响。常见做法是落地到 Hive 前做一次coalesce(1)或repartition(1),把当批数据压缩成少量大文件再写出。代价是写任务变慢,但对后续所有查询都是正向收益。
推荐几个 Hive 侧的合并参数,作为参考:
| 参数 | 含义 | 建议值 |
|---|---|---|
hive.merge.mapfiles | Map-only 任务结束时合并小文件 | true |
hive.merge.size.per.task | 合并后单文件目标大小 | 256000000(约 256MB) |
hive.merge.smallfiles.avgsize | 小于该平均值时触发合并 | 16000000(约 16MB) |
这些参数设好以后,定期对 dwd 层做一次INSERT OVERWRITE重写,能明显改善 Hive 查询的稳定性。这也是我在拆完这个工程后最想强调的一点:Spark 算得快只是起点,Hive 层能不能查得快,靠的是分区设计和文件治理。
5. 避坑指南:Spark+Hive 工程最常见的 5 个翻车现场
5.1 Spark 写完了,Hive 查不到表
现象:Spark 作业正常结束,日志里没有报错,但到 Hive 里SHOW TABLES找不到刚才写入的表,或者表存在但数据为空。
原因:SparkSession 连接到的 Hive metastore 地址和 Hive 客户端连的不是同一个。最常见的是本地调试时 Spark 默认用了内置 Derby 元数据库,而命令行 Hive 连的是远端的 MySQL metastore,两边各写各的,自然互相看不见。
解决:在 SparkConf 里显式指定 metastore 地址,让 Spark 和 Hive 共用一套元数据。
.config("hive.metastore.uris", "thrift://localhost:9083") .config("spark.sql.warehouse.dir", "hdfs://namenode:9000/user/hive/warehouse")这里的关键是 9083 端口必须和hive-site.xml里配置的 metastore 服务一致。只要这个端口不对,Spark 就会默默退回到本地 Derby,两类客户端查到的表完全是两套。
5.2 窗口聚合结果越来越大且充满“重复”
现象:按 10 分钟窗口 5 分钟滑动的逻辑跑了一段时间后,发现同一辆车被统计了很多次,数值和交警发布的断面流量对不上。
原因:滑动窗口天然会重复归属数据。一辆车在 08:03 经过卡口,会同时出现在 [08:00, 08:10) 和 [08:05, 08:15) 两个窗口中。这不算异常,但要分清楚你要的指标口径是“断面流量”还是“窗口累计流量”。
解决:明确口径后选择对应方案。如果报告里定义的是“某时段通过卡口的车辆数”,建议改为approx_count_distinct("vehicle_plate")或者先按车牌在窗口内去重再计数;如果你的指标本身就是“滑动窗口累计过车数”,那就保持原样,并在文档里写清楚口径,不要在答辩时被问住了。
5.3 时间字段全部解析成 NULL
现象:过车流水读进来了,但pass_time这一列在结果表里全是 null,后面所有时间相关的窗口聚合都失效。
原因:设备上报的时间格式不统一。比如有的设备是2024-06-01 08:12:33,有的是2024/06/01 08:12:33,还有的是带时区的 ISO 8601 格式。Spark 在推断 JSON schema 时对混合格式的字符串字段会判定为不可解析,直接置空;或者代码里用了SimpleDateFormat("yyyy-MM-dd HH:mm:ss")去解析,遇到2024/06/01这种格式抛了异常。
解决:在清洗阶段把所有时间字段统一成标准格式,再做后续计算。
from pyspark.sql.functions import regexp_replace, to_timestamp clean_stream = flow_stream \ .withColumn("pass_time_clean", to_timestamp(regexp_replace(col("pass_time"), "/", "-"), "yyyy-MM-dd HH:mm:ss"))这里regexp_replace先把斜杠统一替换成横杠,再交给to_timestamp按固定格式解析。经过这一步,后续的 watermark 和窗口计算才能拿到有效事件时间。我的习惯是任何原始字段进入明细层之前,先做一次“格式验收”,宁可多清洗一步,也不要到聚合阶段才发现时间字段是空的。
5.4 本地跑得好好的,提交到集群就 OOM
现象:本地 IDE 用local[*]模式运行完全正常,打包提交到三台节点的集群上,作业跑几分钟就报 Executor Lost 或 Driver OOM。
原因:本地模式和集群模式的数据分布完全不同。本地模式时数据全在一个 JVM 里,Spark 会投机取巧做简单处理;到了集群,数据被划分到各个 Executor,如果你的代码里有collect()、take()这类把全量数据拉回 Driver 的操作,或者每个 Executor 的默认内存参数没调,就会触发 OOM。
解决:先全局检索代码里有没有collect(),凡是存在的一律改成写出到临时表或 HDFS;然后在提交命令里显式指定内存参数。
spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 4g \ --executor-cores 2 \ --num-executors 3 \ --driver-memory 2g \ target/traffic-teach-1.0.jar--driver-memory 2g对应的是 Driver 端 JVM 堆大小,如果你的 Driver 上确实要做小规模的汇总,这个值可以调到 3g,但不要超过物理内存。Executor 内存 4g 是一个相对保守的起步值,跑通后再根据 Spark UI 里 Storage 页面的实际占用动态调整,不要一上来就写 8g。
5.5 Spark 写 ORC 表时卡在最后阶段
现象:foreachBatch 写 Hive 表的任务一直停留在最后几个 task 上,日志里没有明显的 Exception,但作业就是结束不了。
原因:常见原因是写出时目标目录被占用,或者目标表的小文件太多导致写任务要频繁和 NameNode 通信。还有一种可能是你用了insertInto,但目标表是 TEXTFILE 格式,Spark 需要做一次完整的格式转换,转换过程中的数据倾斜让某一个 Executor 扛了大头。
解决:把insertInto改成先repartition控制输出文件数,再写表;同时确认目标表的存储格式和写入端一致。
batch_df \ .repartition(2) \ .write \ .mode("append") \ .format("orc") \ .partitionBy("dt") \ .saveAsTable("dwd_traffic.traffic_flow_window")这里repartition(2)的作用是把当批数据压成两个文件再落盘,避免一次微批写出几十个几十 KB 的小文件。很多同学的毕设只验证到“能跑出结果”这一步,忽略了文件数对生产查询的影响,如果你能主动控制输出文件数量,这本身就是答辩里一个很加分的技术细节。
6. 进阶:用卡口关系做 OD 分析,让研判系统从“统计”走向“溯源”
统计完各卡口的车流量,整个工程其实还有一张牌可以打,就是过车流水里天然包含的“车辆轨迹”。同一辆车在一个时间段内陆续经过多个卡口,把这些卡口按时间串起来,就能近似还原它的出行路径,也就是 OD 分析(Origin-Destination)。
OD 分析的价值在于:它回答的不是“这个路口车多不多”,而是“这些车从哪来、到哪去”。比如早高峰时段某个区域涌入大量车辆,这些车的前序卡口在哪里,就能判断潮汐车流的主要来向,给信号灯配时提供依据。这里我用 Hive 自连接来实现一个基础版。
SELECT a.vehicle_plate, a.camera_id AS origin_camera, b.camera_id AS dest_camera, a.pass_time AS start_time, b.pass_time AS end_time, (unix_timestamp(b.pass_time) - unix_timestamp(a.pass_time)) AS travel_seconds FROM dwd_traffic.traffic_flow_action a JOIN dwd_traffic.traffic_flow_action b ON a.vehicle_plate = b.vehicle_plate WHERE a.dt = '2024-06-03' AND b.dt = '2024-06-03' AND b.pass_time > a.pass_time AND (unix_timestamp(b.pass_time) - unix_timestamp(a.pass_time)) BETWEEN 60 AND 1800 AND a.camera_id <> b.camera_id;这个自连接的核心约束是最后三行:时间严格递增,间隔限制在 1 到 30 分钟,且两个卡口不能相同。为什么是 60 秒到 1800 秒?间隔太短可能是摄像头重复抓拍同一辆车,不构成有效的卡口间通行;间隔太长则可能是车辆中途停靠,不能代表一次连续的出行。这个区间值的设定其实要按城市道路的实际情况调,如果你处理的是高速路卡口,建议把上限拉到 60 分钟。
查询出来的结果就是一个最简单的 OD 记录:某辆车在某个时刻从 A 卡口出发,在某个时刻到达 B 卡口。如果你想继续聚合出“区域间 OD 量”,再 join 一次摄像头维度表,把 origin_camera 和 dest_camera 替换成区域名,然后以区域对分组计数即可。
Spark 侧也能做同样的逻辑,而且处理大数据量时性能更好。做法是用结构化流或批处理读完一天的数据后,按车牌做窗口内的sortWithinPartitions保证同车记录按时间有序,再通过窗口函数把每条记录的“下一个卡口”找出来,本质上就是上面 SQL 的 DataFrame 版本。这个实现更细,但代码量也更大,课设阶段用 Hive 把链路跑通,再在文档里说明 Spark 化改造思路,已经足够体现你对这条业务线的理解深度。
用这个工程做完 OD 分析后,我才真正理解“交通智能研判”里“研判”两个字的分量:实时车流只是体检报告,OD 路径才是病因分析。从那以后,我每次拿到一个新的 Spark+Hive 相关工程,都会强制自己先看一遍有没有“事实表 + 维度表”这对组合,找到它们再谈调优和扩展。这个习惯帮我避开了很多拿着源码却无从下手的尴尬,希望也能帮到你。
本文还有配套的精品资源,点击获取