ARTICLE DETAIL

资讯详情

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

交通智能调度实战:Spark清洗、Hive特征工程与决策模型全解析

交通智能调度实战:Spark清洗、Hive特征工程与决策模型全解析 前年我接手一个城市交通智能调度项目本来以为算法才是核心真正做进去才发现大数据在交通领域的深层难点是把零散、脏乱的原始数据洗干净再让它真正参与决策。项目开始第一个月GPS轨迹在市区主干道大面积漂移网约车订单和轨迹对不上卡口过车数据有近15%的重复记录调度员看大屏数据时眼神里全是怀疑。这篇就把一个完整的交通智能调度优化链路写出来覆盖集群规划、数据接入、Spark清洗、Hive特征工程、调度决策模型和最终的可视化大屏重点记录踩过的坑和调整过程给正在做交通大数据、网约车数据分析或智慧城市项目的同行一份可以直接参考的实战复盘。1. 交通调度为什么非得上大数据传统经验模型卡在哪1.1 传统调度的问题固定配时与“拍脑袋”派单先说传统交通调度是怎么工作的。绝大多数城市的公交发车间隔是固定时刻表早晚高峰明明客流波动剧烈车辆还是按老节奏跑经常出现某条线路挤到人贴门、司机申请加车调度员还在按上个星期的平均间隔做决定。信号灯配时也是固定周期很多路口晚上十一点车流已经很稀了依然要等90秒红灯。网约车平台早期的派单逻辑更粗暴“谁近派谁”看似合理实际会造成区域运力一边倒城东爆单城西空驶司机在城西苦等二十分钟接不到一单。这些问题的共性在于三个缺陷反馈周期太长空间粒度太粗预测能力为零。所谓反馈周期长指的是调度员看到的是上一个小时甚至昨天的数据空间粒度粗指的是一个片区几百辆车根本看不出哪条路、哪个路口在积压预测能力为零就是只能对现状被动响应没法对未来十五分钟的需求做预判。路况是会传染的晚高峰一个路口堵住二十分钟后相邻三个路口全部瘫痪。没有预测手段调度只能一直追着问题跑。1.2 城市交通数据量的真实体量有人会问交通调度用了这么多年经验为什么突然需要大数据直接看数据量就明白了。按一座中等城市估算全市出租车加网约车约一万五千辆GPS回传间隔三到五秒一辆车一天产生1.7万到2.9万条轨迹点仅浮动车数据一天就有2.5亿条左右。加上卡口过车数据、公交刷卡记录、信号灯检测器流量、气象数据、大型活动事件一天新增的记录量能到几十亿条原始文本落盘大约6到10TB。这个体量Excel肯定没戏传统关系型数据库也不是完全不能处理但调度决策要的是分钟级甚至秒级出结果。比如晚高峰突然下暴雨全城叫车需求十分钟内翻倍系统得马上重新分配运力跑个聚合SQL要半小时调度员早就被乘客电话淹没了。所以HDFS加Spark加Hive这套分布式组合在这个场景下几乎是必然选择HDFS扛存储Spark扛清洗和实时计算Hive扛离线特征加工。1.3 大数据在调度里的角色从“事后统计”到“事前决策”我在项目里对团队强调过一句话大数据在交通调度里不是用来做报表的是用来做决策的。传统报表解决的是“昨天发生了什么”智能调度要回答的是“接下来十五分钟会发生什么我该怎么提前动”。实际落地上核心是两件事。第一件是预测把过去三个月每个十五分钟粒度上的叫车需求、路况、客流数据拿来构造样本让模型学会在什么条件下需求会涨、哪里会堵提前分配运力和调整信号灯方案。第二件是匹配预测之后要把决策落到具体对象上比如哪辆车去接哪个订单、哪条线路需要临时加车、哪个路口需要延长绿灯。这两件事没有大数据基础都做不了预测需要海量历史样本匹配需要实时的全量运力与需求状态而这两类数据恰好都是大数据的强项。2. 集群规划与数据接入从裸机到能跑Spark作业2.1 集群部署策略机器配置、软件栈和资源池划分我们当时是五台服务器起家的集群一台主节点加四台从节点。主节点32核128G内存从节点16核64G内存每台机器挂四块4TB数据盘网卡万兆。这个配置对中型城市的交通数据是够用的前三个月跑下来CPU峰值在60%左右没有出现资源打满的情况。软件栈选型上Hadoop用3.x版本Spark用3.2Hive用3.1ZooKeeper三节点Kafka两台就够。这里有个容易忽略的点YARN资源池一定要分成两条队列离线队列占70%资源实时队列占30%。如果混在一起离线任务一跑起来实时任务全被堵在队列外面调度决策延迟直接飙到十分钟以上。我们后来给实时队列加了容量调度器配置核心参数是capacity30maximum-capacity40保证即使离线任务再多实时计算也有最低资源保障。HDFS副本数我建议设置成3小规模集群数据安全是第一位的别为了省磁盘抠副本数。数据盘必须独立挂载不要和系统盘混在一起否则系统日志写满磁盘NameNode直接进入安全模式整个集群罢工。数据目录规划上我把HDFS数据目录、YARN日志目录、Spark临时目录分别指向不同的磁盘分区遇到问题排查时互不干扰。2.2 数据接入管道Kafka实时流与批量文件双通道交通数据来源分两类接入方式完全不同。第一类是实时流数据包括GPS定位和网约车订单状态。GPS终端通过MQTT或TCP长连接上报后端服务统一写入Kafka按业务拆成独立topicgps_track存轨迹点order_event存订单状态变更traffic_flow存信号灯检测器流量。Kafka分区数设计上要算一下一个分区在普通机械盘上每秒大概能写几千到一万条GPS全城峰值每秒约五千条写入理论上八个分区就够我最后设了十六个分区留出大型活动或极端天气下流量翻倍的余量。分区太少会导致单分区写入超限分区太多又会让下游Spark并发拉取压力变大这个平衡点要结合自己的数据量去调。第二类是批量文件数据比如卡口过车记录、公交刷卡数据很多以五分钟一个文件的形式从第三方系统推送过来先落FTP或对象存储再用DistCp定时拉入HDFS原始目录。这里要注意文件的时间对齐问题卡口数据是按设备本地时间生成的不同设备时钟漂移会造成文件内数据时间戳和文件名时间不一致落地后一定要以数据内容里的时间戳为准重新分区不能直接看文件名时间。2.3 数据分层ODS、DWD、DWS怎么划分数据进集群之后不能直接给模型用我按标准数仓思路分了三层ODS层原始数据照单全收JSON解析成结构化字段直接落表GPS轨迹、订单明细、卡口数据全放这层。这一层保留全部历史作为排查问题的底账。DWD层清洗、去重、修正漂移之后的明细事实表。比如GPS轨迹清洗后的dwd_gps_track_clean订单事实表dwd_order_fact字段口径统一能回答“某辆车某时刻在哪”这种明细问题。DWS层按区域、时间粒度汇总的服务数据。比如每五分钟区域供需比dws_region_supply_demand_5min路段平均速度dws_road_speed_15min这是直接给调度模型和可视化大屏供数的层。分层的好处很实际调度模型不用直接面对脏数据离线任务和实时任务共用DWS数据口径统一。我见过一些项目为了省事模型直接查ODS层结果清洗逻辑改一次模型结果变一次根本没法上线。3. 清洗GPS轨迹和订单数据的那些坑Spark作业调优实录3.1 GPS漂移清洗先用启发式规则别一上来就上HMMGPS漂移是最常见也最头疼的问题。城市峡谷、隧道、地下停车场里终端信号反射导致位置跳变前一秒还在A路口下一秒跳到三公里外的B路段再下一秒又跳回来。如果不处理后面算路段速度、算区域运力全部失真。地图匹配HMM效果好但工程复杂度高我们第一版没有上先用启发式规则把明显漂移过滤掉。核心规则是用前后两个轨迹点的球面距离除以时间差计算瞬时速度如果速度超过给定阈值且连续三个点都超判定为漂移。城市道路正常车速基本不会超过120km/h所以阈值设为120高速公路场景单独放开到150。from pyspark.sql import functions as F def haversine_distance(lat1, lon1, lat2, lon2): # 球面距离计算返回米 ... df df.withColumn(prev_lat, F.lag(lat).over(Window.partitionBy(device_id).orderBy(ts))) df df.withColumn(prev_lon, F.lag(lon).over(Window.partitionBy(device_id).orderBy(ts))) df df.withColumn(prev_ts, F.lag(ts).over(Window.partitionBy(device_id).orderBy(ts))) df df.withColumn( instant_speed, haversine_distance(F.col(lat), F.col(lon), F.col(prev_lat), F.col(prev_lon)) / ((F.col(ts) - F.col(prev_ts)).cast(long) / 1000.0) ) df df.withColumn( is_drift, F.when(F.col(instant_speed) 120, 1).otherwise(0) ) # 连续3个漂移点则过滤整段 df.createOrReplaceTempView(gps_with_flag) cleaned spark.sql( SELECT * FROM ( SELECT *, SUM(is_drift) OVER (PARTITION BY device_id ORDER BY ts ROWS BETWEEN 2 PRECEDING AND CURRENT ROW) AS drift_cnt FROM gps_with_flag ) t WHERE drift_cnt 2 )跑完这版规则GPS数据量减少了约7%抽查结果里明显跳变基本消失。启发式规则的优势是快、可解释、易调参等业务稳了再考虑上地图匹配提高精度。3.2 Spark任务资源调优从跑一个多小时到十几分钟第一版清洗任务全量跑历史数据一个多小时才能跑完根本没法支撑每天增量处理。后来看Spark UI定位问题发现shuffle量异常大数据倾斜明显某些executor处理的数据量是其他的五倍以上。调优做了三件事。第一重新设置executor规格driver内存2Gexecutor内存8G每个executor四个core。第二把spark.sql.shuffle.partitions从默认200调到600让shuffle后的分区更均匀减少单分区数据量。第三开启Spark 3.0之后的AQEAdaptive Query Executionspark.sql.adaptive.enabledtrue让Spark在运行时自动合并小分区、优化join策略。调优前后的对比很直观配置项调优前调优后executor内存2G8Gshuffle.partitions200600AQE关闭开启全量清洗耗时70分钟18分钟单日增量清洗耗时12分钟3分钟这里说个心得不要一上来就堆内存先看Spark UI里每个stage的耗时和shuffle量找到瓶颈再针对性调。有一次我们以为内存不够把executor加到16G结果GC时间反而变长任务更慢了后来才知道是分区数太少导致单分区数据量过大加大分区数就解决了。3.3 订单轨迹匹配与状态码口径统一订单表和轨迹表的关联也踩了坑。一开始直接按订单ID和时间窗口join结果一单匹配出十几条轨迹或者明明有轨迹却匹配不到问题出在订单跨天和设备换绑上。后来统一规则订单表关联轨迹表必须同时满足device_id一致、时间在订单开始和结束时间之间并加上订单状态字段过滤只保留载客状态下的轨迹点用于后续的路径还原。更隐蔽的问题是状态码口径。网约车平台不同业务线的订单状态码定义不一致有的用1代表接单、2代表载客有的用10、20还有的是英文枚举。如果不做映射统一算出来的区域供需比会直接错。我们在DWD层建了一张状态码映射维度表把所有来源的原始状态码统一转为内部标准码再往下游分发。这个动作看起来不起眼但直接影响最终指标的正确性比算法调参重要得多。4. 离线特征工程与Hive优化调度决策的“数据底子”4.1 调度特征表怎么设计用一张主宽表承载80%的查询调度模型和可视化大屏需要用到的特征主要集中在区域和时间维度上。我设计DWS层时没有做成一大堆分散的窄表而是用一张主宽表加少量维度表的结构。主宽表按dt hour region_id分区每一行代表一个区域在一个十五分钟窗口内的完整特征需求侧特征叫车请求量、完成订单量、取消订单量、平均响应时长供给侧特征活跃车辆数、空驶车辆数、平均接驾距离路况特征区域内路段平均速度、拥堵里程占比、排队长度均值外部环境特征天气编码、温度、降雨量、是否节假日、周边POI密度建表时存储格式很关键我用ORC加SNAPPY压缩比纯文本格式查询速度快好几倍磁盘占用也小很多。CREATE TABLE dws.dws_region_supply_demand_15min ( region_id STRING COMMENT 区域编码, hour STRING COMMENT 小时, window_start STRING COMMENT 窗口开始时间, demand_cnt BIGINT COMMENT 需求订单数, finish_cnt BIGINT COMMENT 完成订单数, cancel_cnt BIGINT COMMENT 取消订单数, avg_response_sec DOUBLE COMMENT 平均响应秒数, active_vehicle_cnt BIGINT COMMENT 活跃车辆数, idle_vehicle_cnt BIGINT COMMENT 空驶车辆数, avg_speed_kmh DOUBLE COMMENT 平均速度, congestion_ratio DOUBLE COMMENT 拥堵里程占比, weather_code STRING COMMENT 天气编码, is_holiday INT COMMENT 是否节假日 ) PARTITIONED BY (dt STRING) STORED AS ORC TBLPROPERTIES (orc.compressSNAPPY);宽表的好处是调度模型取特征时只查一张表不需要频繁join查询效率高。代价是每次写入要做大量的汇总计算但这个成本在离线环节是可接受的。如果业务后期加了新的特征维度再单独扩展维度表不要轻易动主宽表结构。4.2 Hive小文件问题写任务一小时读任务一下午项目推进到第三周团队成员反馈Hive查询越来越慢一个简单的count都要跑十几分钟。我一看HDFS目录整个人都麻了——一个分区下躺着几万个小文件每个才几十KB。这是Spark写入Hive时的典型问题shuffle分区过多每个分区写出的数据量太小目录下全是小零碎文件。文件一多NameNode内存压力大任务启动时要扫描的元数据量也大查询自然慢。解决分两步走。第一步是治本控制写入时的分区大小。Spark写入Hive前设置合理的spark.sql.shuffle.partitions并加一个DISTRIBUTE BY让数据按目标分区字段聚类减少跨分区小文件。INSERT OVERWRITE TABLE dws.dws_region_supply_demand_15min PARTITION (dt2024-06-01) SELECT ... FROM dwd.dwd_order_fact DISTRIBUTE BY region_id;第二步是治标对已经产生的小文件做合并。ORC格式的表直接执行ALTER TABLE ... CONCATENATE不需要重写数据几分钟就能把几万个小文件合并成几百个。如果表不是ORC格式就用INSERT OVERWRITE重写一次顺便把格式转成ORC。这里提醒一句小文件问题是Hive性能最大的隐形杀手之一处理完文件合并后同一张表的count查询从十五分钟降到了一分半效果立竿见影。4.3 Join倾斜的排查与治理一个从4小时到20分钟的案例还有一个印象深刻的慢SQL优化案例。一张订单事实表Join区域维度表任务跑了四个小时都没结束明显不正常。我先跑了一个统计SQL按join key分组计数发现其中一个区域key的数据量占全表的40%——热门商圈的数据量是普通区域的几十倍这就是典型的join数据倾斜。处理方法是用两阶段join。先把倾斜key的数据单独拎出来用广播方式和小表维度表直接join因为单个热门区域的数据量即使倾斜也远小于全表其余非倾斜key的数据走正常的shuffle join。最后把两部分结果union all起来。-- 第一阶段倾斜key单独join INSERT INTO temp_result_skew SELECT /* BROADCAST(dim) */ t.*, dim.region_name FROM (SELECT * FROM order_fact WHERE region_id HOT_SPOT) t JOIN dim_region dim ON t.region_id dim.region_id; -- 第二阶段非倾斜key正常join INSERT INTO temp_result_normal SELECT t.*, dim.region_name FROM (SELECT * FROM order_fact WHERE region_id ! HOT_SPOT) t JOIN dim_region dim ON t.region_id dim.region_id;同时开启hive.optimize.skewjointrue作为兜底。经过这个拆分任务从四个多小时降到二十分钟以内之后我把所有关键join任务都检查了一遍凡是涉及热门区域的都不再直接join而是提前加这种拆分逻辑。排查思路可以复制先group by看分布再针对性拆不要盲目加大资源。5. 智能调度决策模型从“经验派”到“数据派”5.1 15分钟粒度需求预测先用GBDT跑稳再谈深度学习调度决策的第一步是预测未来十五分钟每个区域的需求量。我们没有一开始就上LSTM、Transformer这些东西而是先用GBDT模型跑稳因为业务上首先要的不是模型多新而是结果稳定、特征可控、出了问题能解释。特征来自DWS宽表历史同时段需求量、前一小时需求趋势、当前在途车辆数、天气、节假日、周边POI密度。训练样本按时间序列构造用过去九十天的数据训练预测未来十五分钟每个区域的需求区间。GBDT对这类表格数据的拟合能力很强训练速度快特征重要性可以直接输出方便跟业务解释为什么某个区域预测值高——是因为历史同时段高还是因为雨天叠加了晚高峰。后期验证稳定后再考虑用序列模型提升精度也不迟。实际运行中模型每天离线训练一次预测结果写入Redis供在线调度服务读取预测值配合实际值做滚动校验如果预测偏差连续三天超阈值就触发重新训练。5.2 车辆派单的供需匹配把派单问题建模成线性分配预测出需求之后要把空闲车辆分配到需求订单上。这个问题的本质是一个线性分配问题有M个空闲司机N个待服务订单目标是最小化乘客平均等待时间同时控制司机的接驾距离避免为了接一单让司机跑六公里空驶。核心上用的是匈牙利算法的变体。构造一个代价矩阵每个单元格表示司机i接到订单j的代价代价包括预估接驾时间、空驶距离惩罚、区域供需失衡惩罚。如果司机和订单不在同一个网格内且距离超过阈值直接把代价设成无穷大避免出现跨半个城去接人的情况。from scipy.optimize import linear_sum_assignment cost_matrix build_cost_matrix(idle_drivers, pending_orders) row_ind, col_ind linear_sum_assignment(cost_matrix) assignments [(idle_drivers[row_ind[i]], pending_orders[col_ind[i]]) for i in range(len(row_ind))]这个匹配策略上线后平均接驾时间下降约9%同时司机的平均空驶里程没有上升说明匹配不是单纯“就近”而是综合考虑了区域供需。这里要注意的是代价矩阵的规模全城同时有上千个空闲司机和订单时直接调linear_sum_assignment几千乘几千的矩阵性能不够需要先按网格分区做预聚类每个区内部再做匹配减少矩阵规模。5.3 信号灯配时优化用排队长度和流量比替代固定配时常规信号灯配时是固定周期各相位绿灯时长按历史经验分配。我们接入信号灯检测器数据后发现一个路口在一天内不同方向的流量比例变化极大早高峰南向北流量是北向南的两倍晚高峰正好反过来。固定配时在这种路口表现很差排队长的方向绿灯不够用排队短的方向绿灯空放。数据驱动的思路是每个控制周期开始时根据各方向实时排队长度和到达流量动态计算最优周期和绿信比。简化公式是周期时长依据路口总流量和饱和度计算各相位绿灯时间按流量比分配再叠加排队长度修正项。排队超过阈值的相位额外增加绿灯排队为空的方向压缩绿灯。实际改造了三个路口做试点其中一个路口早高峰南进口排队原本经常溢到上游路口调整后排队长度下降18%左右北进口车辆延误增加控制在5%以内整体通行效率是提升的。这个方向项目周期较长涉及路口控制器的协议对接但验证了同一套数据基础设施可以复用到不同调度场景。5.4 实时决策链路模型结果怎么送到调度员手里离线预测是在T-1日晚上跑完产生次日全天的分时预测方案但交通状况随时在变决策链路必须支持实时修正。我们在离线预测的基础上叠加了一整套实时监测Kafka里的GPS和订单流通过Spark Streaming实时计算当前各区域供需比和预测值做偏差比较。偏差超过阈值的区域自动触发调度建议比如向该区域推送空驶车辆、调整公交发车间隔、通知信号灯系统增大某个方向绿信比。实时链路的关键是控制延迟。Kafka消费者如果处理不过来事件堆积等计算结果出来时现场状况早就变了。我们的监控告警是Kafka消费延迟超过三分钟就报警超过五分钟触发自动扩容。调度员的工作台页面五秒刷新一次区域状态模型建议以“可执行任务”的形式推送给调度员而不是自动执行——在交通调度里人的确认环节还是很有必要的尤其涉及跨区域协调时机器建议容易忽略一些文档之外的现实约束。6. 调度效果监控大屏用FlaskECharts把优化结果摆出来6.1 大屏指标设计展示之外还要能指导调度项目做到后半段业务方提出要一个可视化大屏把调度效果展示出来。大屏设计上我坚持一个原则不只是给领导看的漂亮图表更要能辅助调度员日常工作。所以大屏选了三块核心内容左侧是区域运力热力图实时展示各区域供需比中间是地图轨迹图层展示活跃车辆分布和订单流向右侧是时序指标卡展示平均接驾时间、响应时长、优化前后对比折线。指标口径是重点。比如“平均接驾时间”必须明确是从乘客下单到司机接单的时间还是从司机接单到上车的时间两个口径差好几倍。我们和大屏团队成员对了三遍口径再找业务方确认最终统一了所有指标的定义和计算SQL。这里建议所有做数据大屏的人上线前列一份指标口径清单和业务方逐条确认避免大屏上线后被质疑数据不对。6.2 Flask接口与ECharts渲染的落地细节技术实现上后端用Flask前端用ECharts这是做数据大屏最顺手的组合之一。Flask负责从DWS表和Redis读数据封装成JSON接口。实时刷新分成两种不需要秒级变化的数据用前端定时轮询三秒一次需要实时推送的数据用WebSocket比如当前供需比的变化。from flask import Flask, jsonify import redis app Flask(__name__) r redis.Redis(host10.0.0.8, port6379, decode_responsesTrue) app.route(/api/supply_demand) def supply_demand(): data r.hgetall(dws:region_supply_demand) return jsonify(data)ECharts渲染时使用geo地图加scatter图层区域聚合值用visualMap做颜色映射红色代表供不应求绿色代表运力充足。地图JSON使用城市geojson文件坐标要和GPS数据的坐标系保持一致这里坑过我们一次GPS数据是WGS84坐标地图底图用的是GCJ02偏移一两公里热力点全飘到隔壁区。后来统一在数据清洗层做了坐标转换问题消失。6.3 大屏数据正确性验证别让业务方看到错误数字大屏数据最容易翻车的是实时刷新环节。某次上线后业务方反馈“区域红色报警了但实际路口明明很空”排查发现是实时任务重跑导致Redis被写入了一批重复数据。后来我们在写入链路加了一个窗口去重逻辑同一条事件只允许进一次实时聚合并且大屏展示的数据每天凌晨和离线DWS表的批统计做一次对账数字偏差超过1%就告警。另外前端渲染性能也踩过坑。ECharts地图上直接在scatter图层画全量几十万个GPS点直接卡爆浏览器实测五万点以上时帧率掉到20页面操作卡顿。解决方案是网格聚合把全城按1公里网格切分每个网格只展示聚合后的车辆数和供需比点数压缩到两千以内大屏刷新瞬间恢复流畅。这个方案对调度员看区域状态没有任何损失反正他们也不会去数单个车辆。最后一个建议是大屏上线后一定要让调度员实际用起来并根据他们的反馈迭代。我们第二批迭代就是根据调度员意见加的“拥堵趋势预测”面板——他们不满足于看当前状态还想知道接下来半小时哪里会堵。这个需求倒逼我们把离线预测模型的结果同步到大屏上产品的价值也上了一个台阶。回看这个项目最大的体会是算法没有想象中重要数据质量和工程链路才是决定成败的环节。我一度花了三周调派单算法提升不到5%后来回头修了一个GPS漂移过滤逻辑调度响应时间反而下降了12%。对刚入行做交通大数据的人我的建议是先把数据链路打通、指标口径统一再谈模型。这个领域不是靠一个黑盒模型通吃的它需要把数据、算法和业务场景揉在一起最后所有技术都要回答一个问题调度员愿不愿意用你给的建议。
返回列表