ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Spark+Hive交通智能研判系统:从表设计到性能调优全解析

Spark+Hive交通智能研判系统:从表设计到性能调优全解析 简介这套基于Spark与Hive的交通智能研判系统源码包面向毕业设计与课程实训场景聚焦实时交通流量分析与历史数据仓库查询两大任务适合希望掌握分布式计算与数据仓库联用的大数据学习者。压缩包共58个文件主体为42个Java源文件配合9个XML配置、properties环境配置及监控数据样本可用于完成数据清洗、实时统计、结果回写Hive等完整闭环。包体仅953KB工程结构紧凑便于快速部署运行。目前已有147人学习。项目涵盖RDD操作、滑动窗口计算、HQL聚合查询等关键知识点并附有pom.xml与工程配置能支撑课程设计答辩或毕业设计演示帮助读者从代码层面理解交通态势研判的落地实现。1. 交通智能研判系统为什么绕不开 Spark 和 Hive交通智能研判系统这个名字听起来像是个大屏可视化项目但真正落到数据工程层面它解决的是「从海量车辆轨迹和卡口过车记录里算出城市交通到底堵不堵、哪条路事故高发、早高峰比昨天提前了几分钟」这一类问题。这类系统最常见的实现底座就是 Spark 加 HiveHive 负责把几十 TB 的历史轨迹数据管理成规范的表Spark 负责在凌晨那两三个小时里把这些表算成研判指标。如果你接手的是一个打包成 zip 交付的交通智能研判系统打开之后大概率会看到两层东西一层是 Hive 的表定义脚本一层是 Spark SQL 或者 DataFrame 的离线计算任务。这套组合在交通数据场景里落地非常稳因为研判本身不追求秒级实时它更看重「每天准时出结果、结果能回溯、口径能对齐」。适合的人群也明确数仓工程师、用 Spark 做离线分析的数据开发以及刚要把交通数据从 Excel 搬到数仓里的团队。本文就把这套系统的关键环节拆开讲清楚从数据分层、Spark 写 Hive 的 API 选择到小文件治理和倾斜调优全按可复现的路径走一遍。2. 研判系统的数据底座Hive 表怎么设计才扛得住交通数据2.1 交通数据的三种典型形态和分区策略交通智能研判系统的数据源通常有三类。第一类是卡口过车记录每条记录包含车牌、过车时间、卡口编号、车道编号、车牌颜色这类数据的特点是量大但字段固定一天轻松上亿条。第二类是 GPS 轨迹点网约车和公交车的定位数据每秒或每几秒上报一次包含经纬度、速度、方向角这类数据需要清洗掉漂移点才能用。第三类是业务侧的结构化数据比如事故工单、信号灯配时方案、施工占道信息量小但关联价值高。设计 Hive 表时第一类和第二类数据必须按天分区分区字段用dt格式统一成yyyyMMdd。不要用小时分区除非你的集群资源非常充裕。按天分区的好处是分区数量可控一年 365 个Spark 读取时说裁剪简单而且交通研判的大部分指标天然是按天、按周、按月汇总的。卡口过车表建议用 Parquet 格式加 Snappy 压缩查询性能和压缩比比较均衡GPS 轨迹表因为要频繁做经纬度范围过滤同样用 Parquet但排序可以考虑按卡口编号做一次sort by这样同一天的数据在文件内是聚集的后面的 join 能省不少 shuffle。建表时有个容易忽略的细节所有时间字段统一存成 string 类型格式yyyy-MM-dd HH:mm:ss不要在 Hive 里用 timestamp 类型。交通数据的来源系统很杂有的给2024-06-01 08:30:00有的给2024/06/01 08:30:00还有的直接给 Unix 毫秒时间戳。如果建表时用了 timestamp清洗时你就会陷入无休止的格式转换统一存 string 后让 Spark 清洗层去解析和标准化Hive 表只负责存储。这也是做交通数仓最常见的经验。CREATE TABLE dwd_traffic_pass_record ( plate_no STRING COMMENT 车牌号, plate_color STRING COMMENT 车牌颜色: 0蓝 1黄 2绿 3黑, pass_time STRING COMMENT 过车时间 yyyy-MM-dd HH:mm:ss, pass_time_ts BIGINT COMMENT 过车时间Unix时间戳(秒), kam_id STRING COMMENT 卡口编号, lane_no INT COMMENT 车道编号, direction STRING COMMENT 方向: E/W/S/N, speed_kmh DECIMAL(5,1) COMMENT 瞬时速度km/h ) COMMENT 卡口过车明细表 PARTITIONED BY (dt STRING COMMENT 分区 yyyyMMdd) STORED AS PARQUET TBLPROPERTIES (parquet.compressionSNAPPY);这里把pass_time和pass_time_ts同时保留是故意的。pass_time_ts用于窗口函数计算时间差比如判断同一辆车连续通过两个卡口的时长pass_time用于最终展示和关联业务表。DECIMAL(5,1)比FLOAT可靠交通研判要算平均车速和拥堵指数浮点误差在汇总时会被放大。2.2 OD 分析和拥堵指数依赖的中间表设计需求侧交通研判系统最核心的三张中间表是断面流量表、OD 表、路段旅行时间表。断面流量表描述「某个卡口在某个 5 分钟窗口内过了多少车」OD 表描述「车牌从入口卡口 A 到出口卡口 B 的行程」路段旅行时间表描述「某条路在两个卡口之间的通行耗时分布」。这几张表的共同点是需要从明细表做聚合或自关联。比如 OD 表的生成逻辑是对同一辆车按过车时间升序排列找相邻两条过车记录如果时间差在合理区间比如 5 到 180 分钟就认为是一次完整行程。这个逻辑在 Spark 里用Window函数实现。对应的 Hive 中间表要按dt分区同时按kam_id做二级分布。二级分布不是分区分桶Bucket是在建表时用CLUSTERED BY (kam_id) INTO 64 BUCKETS这样后续按卡口关联查流量时能大幅减少 shuffle。OD 表的数据量比明细表小一个量级但关联查询频繁所以存储格式可以用 ORC压缩选 ZLIB查询过滤性能比 Parquet 更优。这里有个选型判断明细表用 Parquet汇总表用 ORC混用两种格式没有任何问题因为 Hive 和 Spark SQL 都能同时读这两种格式。3. Spark 读写 Hive 的工程实现先跑通再优化3.1 Spark SQL 和 Hive 的会话集成方式写交通研判任务时Spark 和 Hive 的关系要理清。Spark 不直接读 Hive 的表数据文件而是通过 Hive Metastore 获取表的元数据schema、分区信息、存储路径然后 Spark 自己直接去 HDFS 上读 Parquet 或 ORC 文件。所以代码里必须在 SparkSession 开启enableHiveSupport()否则spark.sql(select * from dwd_traffic_pass_record)会报Table not found。from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(traffic_judge_daily) \ .config(spark.sql.shuffle.partitions, 200) \ .config(spark.sql.adaptive.enabled, true) \ .config(spark.sql.adaptive.coalescePartitions.enabled, true) \ .config(hive.exec.dynamic.partition, true) \ .config(hive.exec.dynamic.partition.mode, nonstrict) \ .enableHiveSupport() \ .getOrCreate()enableHiveSupport()是核心它会自动读取hive-site.xml里的 Metastore 地址。如果没有这个配置文件需要在代码里显式指定config(spark.sql.catalogImplementation, hive)。shuffle.partitions默认是 200但交通数据量通常上亿200 个分区会导致单个任务处理几百万条记录内存压力极大。建议先设 400 到 600后面根据 Spark UI 里的 Shuffle Read 量再调。读写逻辑可以用spark.sql()直接写 Hive SQL也可以spark.read.table(dwd_traffic_pass_record)读成 DataFrame 再用 API 处理。交通研判这种场景SQL 能表达的聚合逻辑用 SQL 写更直观复杂的状态计算比如 OD 行程识别用 DataFrame API 配合Window更顺手。两者混用没有性能差别Spark 执行引擎是同一个选择标准是代码可维护性。3.2 把 Spark 计算结果写回 Hive动态分区的正确姿势研判结果写回 Hive 时最常用的方式是动态分区写入。比如把每天的卡口 5 分钟流量统计写入结果表分区字段是dt而dt的值是计算出来的。如果用静态分区你得先alter table add partition再 insert非常不方便。动态分区让 Spark 自动根据数据里的分区字段值创建分区。from pyspark.sql import functions as F # 读取当天卡口过车明细 df spark.table(dwd_traffic_pass_record) \ .filter(F.col(dt) F.date_format(F.current_date(), yyyyMMdd)) # 按卡口和5分钟窗口聚合流量 flow_df df.withColumn( window_start, F.from_unixtime( F.floor(F.col(pass_time_ts) / 300) * 300, yyyy-MM-dd HH:mm:ss ) ).groupBy( F.col(kam_id), F.col(window_start), F.col(dt) ).agg( F.count(*).alias(pass_cnt), F.round(F.avg(speed_kmh), 1).alias(avg_speed), F.round(F.percentile_approx(speed_kmh, 0.5), 1).alias(median_speed) ) # 动态分区写回 Hive 结果表 flow_df.write \ .mode(overwrite) \ .partitionBy(dt, window_start) \ .format(hive) \ .saveAsTable(ads_kam_flow_5min, partition_overwrite_modedynamic)动态分区写入有几个关键细节要留意。第一partitionBy里的字段顺序决定了 HDFS 上的目录嵌套层级dt在前、window_start在后这样目录结构是dt20240601/window_start2024-06-01 08:00:00/查询时按天裁剪的代价最小。第二mode(overwrite)配合partition_overwrite_modedynamic只覆盖当天对应的动态分区不会把整张表清空重建这是交通日报任务的标准写法。如果漏掉partition_overwrite_modedynamic参数Spark 3.0 以上版本默认只支持静态覆盖会把整个表的旧数据全部删掉出现「昨天的数据没了」的线上事故。percentile_approx这个函数在交通指标里很有用。平均车速容易被极值带偏比如有一辆车在卡口间飞行速度 180 km/h平均车速就被拉高了中位数车速更贴近真实通行感受。percentile_approx是近似算法误差在 1% 以内对交通研判来说绰绰有余。3.3 Spark 读取 JSON 格式轨迹数据的一个处理技巧交通数据源里GPS 轨迹经常以 JSON 格式从消息队列落盘到 HDFS。Spark 读 JSON 时最常见的问题是 schema 推断。轨迹数据里有的字段是嵌套结构比如location: {lng: 116.4, lat: 39.9}如果让 Spark 自动推断遇到空值多的文件容易推断成不同结构最后各个文件读到的是不同的列。解决方法是把轨迹数据先清洗成统一的 schema再写回 Parquet 格式的 Hive 表。from pyspark.sql.types import StructType, StructField, StringType, DoubleType, LongType json_schema StructType([ StructField(plate_no, StringType(), True), StructField(gps_time, StringType(), True), StructField(lng, DoubleType(), True), StructField(lat, DoubleType(), True), StructField(speed, DoubleType(), True), StructField(angle, DoubleType(), True) ]) # 读取原始 JSON 轨迹文件 trajectory_df spark.read \ .option(multiLine, false) \ .schema(json_schema) \ .json(hdfs://nameservice1/data/gps_traj/dt20240601/*.json) # 过滤漂移点速度大于200km/h直接剔除经纬度超出城市边界剔除 trajectory_df trajectory_df.filter( (F.col(speed) 200) (F.col(lng) 113.5) (F.col(lng) 122.5) (F.col(lat) 20.5) (F.col(lat) 30.5) )给 Spark 显式指定json_schema是一个低成本高收益的优化。JSON 文件没有 schema 约束Spark 推断 schema 时要先扫描一遍全部数据这叫做「schema inference 的额外扫描」会让任务多跑十几分钟。指定 schema 后Spark 直接按既定类型读取省掉第一次扫描同时避免了嵌套结构推断不一致的玄学问题。经纬度范围过滤是交通轨迹清洗的必做步骤GPS 信号在城市峡谷中经常漂移不挡掉的话后续算出来的路段平均速度表会混入大量噪声。4. 从明细到研判核心指标计算的完整链路4.1 路段拥堵指数计算从卡口速度到指数映射交通研判系统里最核心的输出是路段拥堵指数。常见算法是把一条路按相邻卡口切成路段每分钟或者每 5 分钟算一次路段平均速度然后跟这条路的「自由流速度」凌晨车少时的平均速度做对比得出拥堵指数。这个指数再映射成「畅通、缓行、拥堵、严重拥堵」四个等级。实现上分三步。第一步用卡口过车表算平均速度这个在 3.2 节的代码里已完成。第二步把卡口编成路段需要一张路段配置表描述每个路段的卡口序列。第三步用路段配置表和流量表 join按路段分组算指数。# 路段配置表: dwd_road_segment_conf # 包含字段: segment_id, road_name, kam_start, kam_end, free_flow_speed, road_length_km seg_conf_df spark.table(dwd_road_segment_conf) \ .filter(F.col(dt) 20240601) # 流量表按卡口聚合5分钟平均速度 flow_df spark.table(ads_kam_flow_5min) \ .filter(F.col(dt) 20240601) # 关联得到路段级速度 segment_speed_df flow_df.alias(f) \ .join(seg_conf_df.alias(s), onF.col(f.kam_id) F.col(s.kam_start), howinner) \ .select( F.col(s.segment_id), F.col(s.road_name), F.col(f.window_start), F.col(f.avg_speed), F.col(s.free_flow_speed), F.col(s.road_length_km) ) \ .withColumn( congestion_index, F.when(F.col(free_flow_speed) 0, F.round(F.col(free_flow_speed) / F.col(avg_speed), 2) ).otherwise(F.lit(-1)) ) \ .withColumn( level, F.when(F.col(congestion_index) 2.0, 严重拥堵) .when(F.col(congestion_index) 1.5, 拥堵) .when(F.col(congestion_index) 1.2, 缓行) .otherwise(畅通) )拥堵指数用自由流速度 / 平均速度计算指数大于 1.5 说明通行速度明显下降。这个公式有个隐含前提平均速度不能为零。如果统计窗口内某路段没有车通过avg_speed为 nulljoin 出来整行都是 nullcongestion_index会算成 null 而不是-1。所以在聚合后应该先filter(F.col(avg_speed).isNotNull())再算指数。实际交通场景中最容易在这里翻车某条路凌晨没车流量统计为 0平均速度为 null指标口径就乱了。4.2 早高峰识别用 Spark 窗口函数找流量突变点研判系统另一个典型需求是判断早高峰起始时间。常见做法是把一天按 5 分钟切片统计每个切片的流量然后找「流量从平峰上升到峰值 80% 的第一个时间点」。用 Spark 的窗口函数实现最方便。from pyspark.sql import Window # 拉取某区域所有卡口5分钟流量 flow_all spark.table(ads_kam_flow_5min) \ .filter(F.col(dt) 20240601) \ .groupBy(window_start) \ .agg(F.sum(pass_cnt).alias(total_cnt)) # 按时间排序计算滑动窗口内的流量均值用于平滑噪声 window_spec Window.orderBy(window_start).rowsBetween(-11, 0) flow_smooth flow_all.withColumn( smooth_cnt, F.round(F.avg(total_cnt).over(window_spec), 2) ) # 找出当日流量峰值定义早高峰起始 流量达到峰值80%的第一个窗口 peak_cnt flow_smooth.agg(F.max(smooth_cnt)).collect()[0][0] morning_peak_start flow_smooth \ .filter(F.col(smooth_cnt) peak_cnt * 0.8) \ .orderBy(window_start) \ .select(window_start) \ .first()这里的滑动窗口rowsBetween(-11, 0)表示当前行和前 11 行一共 12 个窗口也就是 1 小时5 分钟 × 12的流量均值。用滑动均值再判断峰值可以消除单个 5 分钟窗口的突发流量干扰比如附近停车场出口突然放行了一波车。peak_cnt的计算用了collect()[0][0]在小数据量下没问题如果数据量大更稳的做法是把峰值也通过窗口函数算出来再过滤避免 driver 端拉取数据。4.3 事故高发时段研判关联事故工单和卡口数据事故研判的常见思路是统计事故工单在时间和空间上的分布再和卡口流量关联看哪些时段事故率显著高于均值。这个计算里有一个典型陷阱事故工单表很小关联数据却比工单大几个量级直接用join会在大表上反复扫描。正确做法是先对小表做大过滤再做 broadcast join。具体说事故工单表按事发时间过滤出当天的记录然后用join时指定hint强制 Spark 把小表广播到每个 executor。accident_df spark.table(dwd_accident_order) \ .filter(F.col(dt) 20240601) \ .select(accident_id, kam_id, accident_time, accident_type) # 小表广播避免大表扫描事故表 df_result accident_df.alias(a) \ .join( F.broadcast(seg_conf_df.alias(s)), onF.col(a.kam_id) F.col(s.kam_id), howleft ) \ .join( flow_df.alias(f), on(F.col(a.kam_id) F.col(f.kam_id)) (F.col(a.accident_time).between(F.col(f.window_start), F.col(f.window_start) F.expr(interval 5 minutes))), howleft ) \ .select(a.accident_id, a.accident_type, s.road_name, f.avg_speed, f.pass_cnt)在第 10 到 14 行事故时间落在哪个 5 分钟窗口内决定了关联的流量值。如果事故发生在 8:03流量窗口是 8:00 或 8:05between条件用窗口开始时间和上界区间能精确匹配到 8:00 窗口。但如果事故发生在 8:04而流量窗口只有 8:00 和 8:05 两个点between会同时匹配到这两个窗口产生一对多的膨胀关联。这里建议先做窗口对齐把事故时间floor到窗口开始时间再去关联而不是直接between。这是交通数据关联里最隐蔽的性能和口径坑。5. 交通数据任务常见问题与避坑实录5.1 小文件问题从 Parquet 文件块数太多到 HDFS NameNode 告警现象每天跑完日批任务后HDFS 上某个结果表的分区目录下有成百上千个小文件每个只有几十 KB 到几 MB。整张表文件数上百万NameNode 内存告警后续 Spark 查询时任务数暴增单个任务处理的数据量又太少整个集群像在空转。 原因Spark 写入时每个 partition 对应一个文件而shuffle.partitions设得过大比如默认的 200加上动态分区写入时每个动态分区又会生成一批文件结果就是文件数 分区数 × 动态分区数。交通数据的动态分区一天能有上百个按小时或窗口分区时更是上千文件数膨胀得很快。 解决在写 Hive 表之前用repartition或coalesce控制输出文件数。如果结果表的目标是每 5 分钟一个窗口每窗口数据量不大就按方式一处理——先按窗口聚合再repartition(1)让每个窗口只出一个文件如果数据量大控制每个文件 128 MB 到 256 MB按分区键做repartition后再写。# 控制动态分区写入的文件数 flow_df_for_write flow_df.repartition( F.col(dt), F.col(window_start) ) # 按分区键 shuffle保证同一个分区数据到一个 executor flow_df_for_write.write \ .mode(overwrite) \ .partitionBy(dt, window_start) \ .option(maxRecordsPerFile, 2000000) \ .format(hive) \ .saveAsTable(ads_kam_flow_5min, partition_overwrite_modedynamic)这里repartition按分区键进行保证同一个(dt, window_start)组合的数据都在同一个 executor 上写出的文件就是「一个分区一个文件」。maxRecordsPerFile是保险丝防止某些分区数据量特别大时单文件过大。注意不要用coalesce代替repartitioncoalesce只减少分区数不做数据重分布分区键相同的记录可能散落在多个文件里。另外要定期对 Hive 表做小文件合并。常见做法是用 Spark 读全表后repartition写一遍但更实用的方案是直接对历史小分区目录做INSERT OVERWRITE ... SELECT ...只重塑当天分区。5.2 数据倾斜车牌分布不均引发的长尾任务现象ETL 跑卡口流量统计时大部分任务几十秒跑完但有一个任务跑了 40 分钟还没结束Spark UI 里看到某个 stage 的某个 task 处理的数据量是其他 task 的几十倍。仔细观察这个 task 处理的 key 通常是市中心的某个核心卡口比如「人民广场」卡口一天过车 20 万辆而郊区卡口一天只有几千辆。 原因group bykam_id时热点卡口的数据集中在少数几个 partition 上负责这些 partition 的 task 负载极高既拖慢整体任务还可能因为数据量过大触发内存不足OOM。 解决给热点 key 加盐盐值打散。原理是把热点 key 拆分成多个虚拟 key分别聚合后再合并结果。交通场景下热点卡口数量少但流量极大加盐的收益非常明显。# 卡口流量聚合加盐处理 salt_range 10 # 加盐数量热点卡口拆成10个虚拟key flow_with_salt df.withColumn( salt_kam_id, F.when( F.col(kam_id) KA0001, # 识别热点卡口 F.concat(F.col(kam_id), F.lit(_), F.lit(F.floor(F.rand() * salt_range))) ).otherwise(F.col(kam_id)) ) agg_result flow_with_salt.groupBy(salt_kam_id, window_start) \ .agg(F.count(*).alias(pass_cnt)) \ .withColumn( kam_id, F.split(F.col(salt_kam_id), _)[0] # 还原真实卡口编号 ) \ .groupBy(kam_id, window_start) \ .agg(F.sum(pass_cnt).alias(pass_cnt))第一次groupBy在加盐 key 上并行聚合第二次groupBy把虚拟 key 的结果合并还原。注意F.lit(F.floor(F.rand() * salt_range))中rand()和lit的组合如果直接写在when里Spark 会对每一行都计算一次rand()不会取固定值这反而能更均匀地打散数据。生产环境不要硬编码热点卡口而是先从明细表里统计出 Top N 卡口动态生成热点名单。5.3 动态分区写入覆盖了不该覆盖的数据现象某天重跑前一天的研判任务原本只想更新dt20240601的数据结果整个ads_kam_flow_5min表的数据全部消失只剩今天跑的数据。检查代码发现mode(overwrite)把表清空重建了。 原因Spark 3.0 起partitionOverwriteMode默认值是STATICINSERT OVERWRITE TABLE会先删除整张表数据再写入新数据。这跟你是不是只在分区条件里指定了特定日期的值没关系默认静态模式就是全表覆盖。 解决在 SparkSession 里显式设置partitionOverwriteModeDYNAMIC或者在写入时用上面代码里的.option(partition_overwrite_mode, dynamic)明确指定。这个配置必须在写操作前设置写在spark.sql()的 set 里也行。另外如果写的是分区表建议再开启spark.sql.sources.partitionOverwriteModedynamic双保险。这是交通数仓日报任务最容易翻车的点没有之一。5.4 Hive 表分区里出现乱码或日期格式不一致的分区值现象SHOW PARTITIONS dwd_traffic_pass_record输出的分区列表里混着dt20240601和dt2024-06-01这两种格式导致WHERE dt20240601查不到前一天 80% 的数据。 原因不同批次的写入任务用了不同的日期格式化函数有的用date_format(current_date(), yyyyMMdd)有的直接substring(regexp_replace(current_date(), -, ), 1, 8)。如果代码里把yyyy-MM-dd字符串误当成动态分区键的值传给saveAsTable分区目录就变成dt2024-06-01这种形态。 解决在数仓规范层面做强制约束全链路统一用yyyyMMdd并在调度平台给所有日报任务配置一个公共的日期参数变量bizdate统一取到的是格式化好的字符串。在代码层面如果是已有的脏分区需要先清理掉# 在 Hive 中删除乱码分区 hive -e ALTER TABLE dwd_traffic_pass_record DROP IF EXISTS PARTITION (dt2024-06-01);同时把 Spark 写入前的过滤条件改为F.col(dt) bizdate.replace(-, )来规避日期字符串混用。如果脏分区已经非常多直接用MSCK REPAIR TABLE table_name把新识别出的分区加载进来再配合脚本把格式错误目录从 HDFS 上清理掉。5.5 Spark 内存调优执行器内存与并发度怎么配现象凌晨跑交通研判日批任务时任务经常在某个 stage 失败报Container killed by YARN for exceeding memory limits或者出现频繁的 GC 暂停。 原因交通数据groupBy、join的 shuffle 量非常大默认的 Spark 执行器内存1G 到 2G完全不够用。YARN 容器内存达到上限被杀掉时Spark 重试机制会重复拉取数据整个任务越跑越慢直到失败。 解决按数据量估算内存和并发度。经验值是数据量 500 GB 到 1 TB 的日批任务executor 数量 50 到 80 个每个 executor 内存 8 G 到 16 Gspark.sql.shuffle.partitions设为 executor 数的 2 到 4 倍然后交给 Spark 动态资源分配去控制。spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 12g \ --executor-cores 4 \ --num-executors 60 \ --conf spark.sql.shuffle.partitions240 \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.adaptive.coalescePartitions.enabledtrue \ --conf spark.shuffle.service.enabledtrue \ --conf spark.memory.offHeap.enabledtrue \ --conf spark.memory.offHeap.size4g \ traffic_judge.pyspark.shuffle.service.enabled开启后executor 被 YARN 回收时 shuffle 文件不会丢失能避免动态资源缩容引发的宽依赖报错。spark.memory.offHeap.size主要给堆外的 Netty 和加密用对交通数据这种大量 decimal 计算有一定帮助。内存调优没有银弹必须结合 Spark UI 的Shuffle Read Size / Record和GC Time两个指标来判断瓶颈。GC 时间超过总任务时间的 10%就加大spark.executor.memoryOverhead默认是 executor 内存的 10%对于大 shuffle 任务建议调到 15% 到 20%。5.6 小文件合并时机和调参技巧现象跑完两周的日报后某个结果的表数据量没怎么涨磁盘占用却翻了三倍。查看文件发现一个 128 MB 都不到的 Parquet 文件占了 30 MB 的元数据开销或者文件内部被切成了过多的小 RowGroup读取时要解压大量无效数据。 原因动态分区写入 高频插入每 5 分钟调度一次写入每次写入都会在同一个分区目录下追加新文件。日积月累文件数多Parquet 文件的 RowGroup 切成很多小块。Parquet 读取时是按 RowGroup 和 ColumnChunk 为单位做谓词下推的文件碎片化后索引的效果大打折扣。 解决在调度上用「先写临时分区再合并到正式分区」的两阶段模式。具体实现是每天的任务先写入当天的临时分区分号或tmp_前缀写完后用INSERT OVERWRITE一个 SQL从临时分区SELECT出来再重新写到正式分区目录这一步让 Spark 重新生成文件。更趋向于治理的方案是单独跑一个小文件合并任务把过去 N 天的小分区统一重写一遍。-- 先把当天临时分区合并到正式分区 INSERT OVERWRITE TABLE ads_kam_flow_5min PARTITION (dt20240601) SELECT kam_id, window_start, pass_cnt, avg_speed, median_speed FROM ads_kam_flow_5min_tmp WHERE dt 20240601;这个 SQL 的关键是OVERWRITE模式的PARTITION指定了具体分区Spark 只覆盖目标分区不会动其他分区。同时因为 SELECT 是完整的一张表Spark 写出的文件数量和 shuffle partitions 保持一致可以通过SET spark.sql.shuffle.partitions32来控制这次写出的文件数。如果还有残留小文件再用ALTER TABLE ... CONCATENATE对 ORC 表做轻量合并Parquet 表则没有这个命令只能通过重写。5.7 数据不对账研判指标和原始卡口数据对不上现象前端大屏上显示的某路段早高峰平均车速 45 km/h但直接查原始卡口表按公式手算得到的却是 51 km/h。开发人员重新跑了一遍任务结果和第一次跑出来的又不一致每天输出都有细微波动。 原因最典型的两个坑。一是 Spark 的repartition和动态分区写入过程中数据发生了部分重算或重复读取比如读数据时没有做顺延分区读到了dt 昨天的数据或者指标任务和数据抽取任务并行跑数据还没准备好就开始算。二是avg_speed在聚合时包含了瞬时速度为 null 的记录count(*)统计的车辆数和原始明细记录数不一致。 解决在任务代码里对指标结果做「三数对账」原始过车明细表的计数、媒体计数count(*)、去重车牌计数count(distinct plate_no)三个数字定在同一个时间窗口内应满足明确的数量关系。把对账逻辑写成 Spark DataFrame 的assert语句跑批结束后自动校验数字对不上就提示告警。同时所有写 Hive 表之前源表的分区应该等数据完全就绪再开放下游读取这个可以通过调度平台的上游依赖来实现。6. 进阶研判结果回写 MySQL 与增量准实时验证交通研判系统跑完离线任务之后结果数据最终要服务于大屏和移动端。Hive 表不适合直接面对高并发查询这里常见的做法是把 ADS 层结果表通过 Spark 同步到 MySQL让后端服务用 JDBC 查询。Spark 写 MySQL 时有两个参数最值得关注。# 从 Hive 读取当日研判结果 ads_df spark.table(ads_congestion_index) \ .filter(F.col(dt) bizdate) # 写回 MySQL供大屏后端服务查询 ads_df.write \ .mode(overwrite) \ .jdbc( urljdbc:mysql://rm-xxx.mysql.rds.aliyuncs.com:3306/traffic_judge, tableads_congestion_index_daily, properties{ user: data_writer, password: your_password, driver: com.mysql.cj.jdbc.Driver, rewriteBatchedStatements: true, batchsize: 5000 } )rewriteBatchedStatementstrue这个参数非常关键把单条 insert 改成多值批量 insert写 MySQL 的吞吐能提升好几倍。另外一个坑是mode(overwrite)配合 MySQL 时Spark 会把原表DROP TABLE再重建造成大屏查询瞬间报错。稳妥做法是先写一张临时表ads_congestion_index_daily_tmp写入成功后用ALTER TABLE RENAME做原子替换这样前端大屏永远查得到完整数据。对于研判结果的验证最靠谱的方法是把当天结果和前一天结果做对比计算指标波动率超过阈值就告警。比如某路段拥堵指数从 1.2 跳到 2.5要么是真实发生了交通事故要么是数据管道出问题了。这个对比逻辑非常适合用 Spark 写成一个定时校验任务today_df spark.table(ads_congestion_index).filter(F.col(dt) today) yesterday_df spark.table(ads_congestion_index).filter(F.col(dt) yesterday) joined today_df.alias(t).join( yesterday_df.alias(y), on[segment_id], howfull ).select( F.col(t.segment_id), F.col(t.congestion_index).alias(idx_today), F.col(y.congestion_index).alias(idx_yesterday) ).withColumn( change_ratio, F.abs(F.col(idx_today) - F.col(idx_yesterday)) / F.col(idx_yesterday) ) abnormal_segments joined.filter(F.col(change_ratio) 0.5)我个人做交通研判项目最大的教训是不要只信聚合结果。每次上线新指标必须拿 SQL 直接查原始明细表手动算两个窗口的数据去对账。因为 Spark 任务一旦写复杂数据血缘就容易丢失什么时候多算了一个分区、什么时候聚合时漏掉了空值光看输出结果完全看不出来。交通数据链路长、字段杂今天的坑可能埋在一个月前的方言函数里。一个好习惯是保留所有中间层结果表的血缘脚本清洗、聚合、关联各存一份带版本的 SQL这样哪怕线上数据出了偏差也能快速定位是哪个环节引入了脏数据。希望这套从表设计到调优再到对账的路径能帮你在做交通研判和类似的大数据离线分析项目时少走几次弯路。本文还有配套的精品资源点击获取
返回列表