
我最初决定参加天池地铁流量预测这个赛题时真的想简单了。我以为要做的事情很明确——基于历史进出站刷卡数据预测未来某个时段各个站点的客流量这不就是个时间序列回归问题吗直到我尝试把全量数据塞进单机内存做特征工程才意识到城市级别的 AI compute 远没有想象中那么轻松。几千万条带时间戳的地铁客流记录加上站点属性、天气、节假日等外部数据单机方案直接 OOM。也正是从那一刻起我开始认真搭建基于 Distributed-Platform分布式计算平台和 Distributed-Databases分布式数据库的完整处理链路。这篇文章就从我参加天池地铁流量预测的完整过程出发聊聊分布式架构在城市流量预测项目里到底怎么落地、有哪些坑、以及最终如何把 WAPE 一步步降下来。如果你是做时序预测的或者正在纠结要不要给数据处理上大数据框架再或者准备打这类数据竞赛这篇文章应该能给你不少可以直接拿去用的经验。1. 城市级流量预测的算力瓶颈为什么一开始就选择分布式路线1.1 单机方案的真实天花板我最初拿到数据时的第一个想法是直接读进 Pandas靠单机算完所有特征。当时手头的训练集大概有 7000 万条记录覆盖 320 多个站点时间粒度是 15 分钟一条。单独看训练集的存储大小几个 GB 似乎不算夸张但真正开始做特征工程后才发现问题远不止“内存够不够”这么简单。Pandas 加载全量 DataFrame 大约占 8~10GB 内存这已经接近我笔记本的上限了。紧接着做第一个稍微复杂的操作——按站点分组、按时间排序、计算过去 7 天同一时段的历史均值——中间过程产生的临时 DataFrame 又吃掉了好几个 GB。结果就是内存报警进程被系统杀掉。就算换了内存更大的机器每次迭代特征都要把整个流程重跑一遍一次实验动辄一两个小时一天根本做不了几轮。这里有一个很关键的技术认知Pandas 的底层操作虽然大量使用了向量化指令但在处理 groupby 之后的复杂跨组逻辑时性能并不理想而且它是单进程模型用不满多核 CPU。也就是说数据规模一旦跨过某个阈值单机方案的瓶颈本质上是“内存带宽 CPU 核心数 迭代效率”三重叠加不是单纯靠加内存就能解决的。分布式平台的核心思路是分而治之——把数据按照站点、时间或者其他业务维度切分到多个节点上每台机器只处理自己负责的那一摊最后再合并结果。它的价值不仅仅在于“能放下更大的数据”更在于“同样的特征计算任务可以在十几分钟内完成”这才给了算法工程师高频迭代的可能性。1.2 平台与数据库选型Spark ClickHouse 的组合逻辑我用了 Spark 作为分布式计算平台ClickHouse 作为分布式分析型数据库。这个组合并不是一上来就定死的中间也纠结过 Flink、HBase、MongoDB 这些方案。最终定下来的理由很简单这个赛题的数据是 T1 批次提供的不是实时流Flink 主打的毫秒级低延迟在这里没有用武之地特征工程里 90% 的操作是“按时间范围 站点维度做聚合统计”这是 ClickHouse 这类 OLAP 数据库的强项Spark 的 DataFrame/SQL API 在处理“宽表拼接、多源 join、复杂窗口逻辑”时足够灵活生态成熟遇到问题时社区案例多。我也在架构里预留了一条 Flink 的通道但那是为了事后模拟线上实时预测场景用的。比赛主流程Spark 完全够用。分布式数据库的选型对比我整理了一个简表这对我后面做技术决策帮助很大数据库类型优点缺点本次的定位ClickHouseOLAP 列式存储聚合查询极快支持 SQL列式压缩节省空间点查弱不适合高频单行更新主力明细存储 特征预计算HBaseKV / 宽表海量数据点查能力强支持实时读写聚合计算弱运维复杂备用方案最终未采用Redis内存 KV读写延迟极低容量有限不适合存全量只作为特征缓存辅助MySQL关系型事务能力强生态成熟数据量大后查询慢只存站点元数据、节假日表最终选择 ClickHouse 还有一个务实的原因它支持标准的 SQL 语法团队里哪怕不熟悉大数据技术栈的同事也能很快上手写聚合查询。分布式系统最怕的不是性能不够而是技术门槛高到没人愿意维护。1.3 分布式数据库的表结构与分区策略这是从 MySQL 思维切到 ClickHouse 时最容易踩坑的地方。我建的客流事实表结构大概是这样的CREATE TABLE metro_flow ( station_id UInt32, datetime DateTime, in_count UInt32, out_count UInt32, date Date MATERIALIZED toDate(datetime) ) ENGINE MergeTree PARTITION BY toYYYYMM(date) ORDER BY (station_id, datetime)这里的两个细节非常关键。第一个是分区键。我按月份分区因为特征计算大多只需要最近 28~60 天的数据按月份分区可以快速跳过不需要的月份减少扫描量。这也是时序数据最自然的组织方式。第二个是排序键的顺序。我一开始写成了ORDER BY (datetime, station_id)结果查询某个站点一个月的数据时ClickHouse 几乎做了全表扫描。原因在于 MergeTree 的稀疏索引是严格按照排序键建立的datetime放前面意味着同一个站点的数据在物理存储上被时间戳打散了。改成(station_id, datetime)后单个站点的数据在磁盘上连续存放查询效率提升了大概一个数量级。分布式数据库的存储设计直接决定了上层 AI 计算能跑多快。表结构没设计好后面所有特征计算都会卡在 IO 上换更快的 CPU 也没用。这是我整个项目里体会最深的一句话。2. 数据管道拆解从官方CSV到高效特征宽表2.1 端到端架构与组件分工不看架构直接写代码是打比赛的大忌。我最终跑通的链路大致是这样的官方 CSV 文件先用 DataX 批量导入 HDFS从 HDFS 通过写入任务落到 ClickHouse 明细表Spark 通过 JDBC 读取 ClickHouse 原始数据跑全量特征工程生成的宽表特征写回 ClickHouse同时导出 Parquet 用于模型训练LightGBM 读取特征宽表训练模型预测阶段对测试时间点逐站逐时段生成特征产出最终结果。这套链路里有一个非常重要的原则Spark 和 ClickHouse 的分工边界必须清楚不能重叠。ClickHouse 负责“快查快算”比如按天聚合出每个站点每天的总客流量这类简单聚合直接下推给 SQL 执行秒级返回Spark 负责“复杂特征和宽表拼接”比如多表 join、自定义函数、跨时间窗口的复杂特征逻辑。如果把所有计算都拉到 Spark 里做集群压力会很大而且 ClickHouse 的列式聚合优势就白白浪费了。我实测下来的感受是将简单的聚合下推给 ClickHouse 后Spark 任务的耗时大约下降了 30% 左右。这种“把合适的计算放在合适的引擎里”的思路比单纯优化某一个框架的并发参数要有效得多。2.2 ETL 中最容易忽略的两个细节时间类型与增量设计第一个细节是时间字段的转换。官方给的时间字段是字符串如果直接导入 ClickHouse会被当作字符串排序按时间范围查询时根本无法利用索引。我在导入前统一转成 DateTime 类型并且在并行导入时保证数据按时间戳严格递增写入。这样 MergeTree 内部的每一个 data part 都是有序的后续查询性能会好很多。第二个细节是增量与全量的取舍。这个赛题的数据是一次性提供的理论上是全量但我还是按天做了增量接口的设计。原因是在模型验证阶段我需要模拟线上效果——用前 N 天训练用第 N1 天验证。如果每次回溯验证都重新全量加载时间成本完全不可接受。增量接口让回溯测试变得非常顺畅。我封装了一个函数传入起始日期和结束日期自动从 ClickHouse 读取对应时间范围的数据生成特征宽表并返回训练样本集。这个设计在后面跑滚动验证的时候帮了大忙。2.3 数据一致性分布式环境下的隐形杀手分布式环境下最容易被忽视的是数据一致性。我踩过一个印象很深的坑用 Spark 通过 JDBC 并行读取 ClickHouse 数据时每个分区独立建立连接、独立查询由于多个副本节点的数据同步存在微小延迟某些分区读到的快照在时间边界上并不完全一致。结果就是特征值出现莫名其妙的“跳变”同一站点同一时段在某些实验里特征值不同模型训练出来的结果也忽高忽低。后来我规定所有 Spark 读取任务的 SQL 里都必须带上一个固定的时间边界条件例如WHERE datetime 2023-01-01 00:00:00 AND datetime 2023-02-01 00:00:00确保所有 executor 读取的都是同一个时间切面的数据而不是依赖默认的并行分区切割逻辑。另外一个相关的问题是分布式导入时的数据重复。由于网络重试或者其他原因同一条记录可能被写入两次。在单机 MySQL 里主键约束可以拦截重复写入但 ClickHouse 的 MergeTree 默认没有这个能力。我的解决方案是改用 ReplacingMergeTree 引擎并按(station_id, datetime)设置去重键配合版本列保留最新记录。这些看似和数据预测无关的细节恰恰是分布式系统里最影响模型稳定性的地方。特征不一致导致的模型抖动比模型本身调参不当更难排查。3. Spark 上的特征工程高效写法与拖垮集群的反模式3.1 四类特征的完整设计这个赛题的核心任务是预测某站某时段进出站人数我最终把特征分成四大类总共约 130 维时间特征小时、分钟、星期几、是否工作日、第几周、是否节假日、节假日前后偏移天数。这些特征直接编码了客流周期性规律。历史统计特征过去 1/3/7/28 天同时段的均值、过去 7 天最大值、上周同日同时段值、昨日同时段值。这类特征的预测力最强。站点属性特征站点日均客流量级、是否为换乘站、站点周边土地性质办公/居住/商圈。这些特征帮助模型区分不同类型站点的客流模式。外源特征天气温度、降水、节假日安排。城市地铁客流受天气和节假日影响非常大尤其是一些地面通勤占比较高的线路。从重要性排序来看历史统计特征贡献了绝大部分预测力时间特征次之外源特征在特定日期如节假日前后、暴雨天能起到关键修正作用。3.2 滑动窗口特征的两种实现与性能对比滑动窗口特征是时序预测最有用的特征但在 Spark 里实现不好会非常慢。我最初的写法是使用窗口函数from pyspark.sql import Window from pyspark.sql import functions as F window_spec ( Window .partitionBy(station_id) .orderBy(datetime) .rowsBetween(-28 * 96, 0) # 28天每天96个15分钟时段 ) df df.withColumn(mean_28d_same_slot, F.avg(in_count).over(window_spec))这段逻辑在数据量几百万的时候没问题但到 7000 万行时跑得非常痛苦。partitionBy orderBy会引发全量重分区而且窗口内部的全局排序是一个非常重的操作单个 stage 跑了快一个小时都没结束集群资源全部被占满。后来我换成了“维度逐级提升”的思路先用 ClickHouse 或者 Spark SQL 按“日 站点 时段”聚合出日粒度表在日粒度小表上做滚动窗口聚合计算过去 28 天同时段均值把结果回填到原始的 15 分钟粒度数据上。日粒度表的数据量只有原始表的几十分之一窗口函数跑起来毫无压力。这种“先降维聚合再回填扩展”的做法适用于几乎所有滚动统计特征能把计算量降一个数量级。最终全量特征工程从原来预计的一小时以上压缩到了 25 分钟左右。3.3 外源特征 join广播变量是唯一正确选择吗天气特征、节假日特征这些外源数据通常很小几百行到几千行但如果不注意就会和主表发生 shuffle join白白增加网络传输开销。正确的做法是使用 Spark 的 broadcast join把这些小表广播到每个 executor 上直接在内存里完成匹配几乎不增加集群负担。这里有一个容易被忽略的点并不是所有小表都适合广播。如果外源表有几十万行广播反而会占用大量 executor 内存拖慢整体性能。我的判断标准是单表大小在几十 MB 以下才适合广播。对于更大的外部表先用 SQL 在 ClickHouse 里完成预聚合再以较小粒度与主表关联。地铁站点的 POI 数据、线路属性这类静态数据我也都打成了广播变量。实测下来这类优化对整体任务耗时的改善非常明显尤其是多次迭代实验时省下的时间会累积成很大的优势。4. 模型与训练LightGBM 在分布式环境下的实战调优4.1 为什么最终选了 LightGBM 而不是 LSTM很多人听到时序预测第一反应就是上 LSTM。我也专门花了一周时间搭了一个两层 LSTM 的实验效果却差强人意。相比之下LightGBM 在这个任务上全面胜出。原因并不复杂。地铁客流数据具有极强的周期性和节假日效应星期几、是否节假日、过去几周同时段的客流均值这些特征几乎已经“锁定”了预测值的核心波动范围。树模型对这类离散周期特征的拟合能力远胜于深度学习模型而且不需要像 LSTM 那样精细地调整序列长度、隐层维度、dropout 等一大堆超参数。还有一个很现实的因素是迭代效率。LSTM 在分布式环境下训练需要额外的参数服务器架构训练一轮就要几十分钟调参周期被拉得很长。而 LightGBM 的直方图算法天然适合分布式训练过程中每个 worker 先在本地构建直方图再通过 allreduce 做全局合并通信量小、收敛快。我在 4 台 8 核 16GB 的机器上训练全量特征大约 40 分钟就完成了。4.2 统一建模、分站点建模还是分组建模我在实验过程中对比了三种策略统一建模所有站点放一起训练模型规模大泛化性好但对部分长期客流偏低的站点容易产生系统性偏差。分站点建模每个站点训练一个独立模型拟合能力强但部分站点历史数据太少小模型容易过拟合而且 320 个模型的训练和线上维护成本都很高。分组统一建模把站点按客流量级别、站点周边属性分成若干组组内站点共享一个模型。实测下来分组统一建模的线上效果最佳比纯统一模型高大约 1~2 个百分点。道理也简单大型换乘站和社区型小站的客流模式完全不同硬塞进同一个模型里模型只能学到一个“折中”的模式而分组之后模型能学到更具针对性的客流规律。分组数量我最终定为 6 组每组大约 50~60 个站点。这个“中间态”思路也适用于很多生产场景——既保留统一模型的鲁棒性又兼顾不同业务的个性化差异关键在于找到合适的分组维度。4.3 分布式训练的参数细节与资源规划我用的是 LightGBM 的分布式训练接口最终选择基于 Dask 的实现核心参数如下import lightgbm as lgb params { objective: regression, metric: mae, num_leaves: 127, learning_rate: 0.03, feature_fraction: 0.8, bagging_fraction: 0.8, bagging_freq: 1, lambda_l2: 1.0, min_data_in_leaf: 100, max_bin: 255, num_boost_round: 3000, early_stopping_round: 100, }一个很实际的调优经验是在分布式环境下不要盲目增大数据量去换取精度的微小提升。数据量越大直方图合并的网络通信开销就越大训练时间会非线性增长。当测试发现数据量超过一定规模后评估指标几乎不再变化时就应该停下来而不是继续加数据。另一个经验是特征类型预处理。LightGBM 默认不开启类别特征处理因此星期几、站点 ID 这类特征都需要手动编码。对于站点 ID 这种高基类别特征直接用 label encoding 很容易让树模型在无意义的高基数 ID 上过拟合。我换成“目标编码target encoding 频率编码”的组合也就是用该站点历史平均客流量以及该站点在总样本中的出现频率来替代原始 ID。这样既保留了站点的个性化信息又避免了高基数特征带来的过拟合风险。5. 调优实录WAPE 从17%到12%的改进路径5.1 读懂 WAPE它和 MAE、RMSE 有什么不同天池这个赛题使用的是加权绝对百分比误差WAPE作为核心评估指标公式是WAPE Σ|actual - forecast| / Σactual这个指标的直观含义是“总体预测误差占总体真实客流量的比例”。它和 MAE、RMSE 最大的区别在于它不是对所有样本点误差取平均而是先把所有样本的真实值相加作为分母再把所有样本的绝对误差相加作为分子。这带来的影响是——客流量大的站点在指标中天然占据更高的权重。这在业务上是合理的因为调度优化的重点本来就应该放在人流量大的站点单个小站在数值上影响很小。明白了这一层调优方向就清晰了与其在几百个小站点上花精力把误差从 50% 降到 40%不如在几个核心大站上把误差从 10% 降到 8%后者对 WAPE 的贡献更大。5.2 高峰与节假日的漂移三个典型翻车场景我的调优过程分为几个阶段每次改进都对应一个具体的翻车场景。第一个场景是节假日前后一天。节假日前一天晚上和节后第一天的客流模式完全不同于普通工作日过去几周的历史均值特征在这里反而会成为误导。我的解决方案是增加“节假日前后偏移天数”特征把节前 1 天、节中、节后 1 天分别标记为不同的状态。这个特征加入后节假日样本的预测误差明显缩小。第二个场景是极端天气。大雨天地铁客流量会比平时有所增加尤其是地面公交出行受影响时但具体增幅在不同站点差异很大。应对方式是把天气特征从“当天的天气概况”细化为“过去 3 小时降水量和温度变化”让模型能捕捉到天气变化的短期效应。第三个场景是周一早高峰和周五晚高峰。这两个时段与普通工作日有明显差异。我的处理办法是把“是否周一早高峰” “是否周五晚高峰”作为独立特征加入模型而不是让模型从“星期几 小时”的组合里自己去学。经过这一轮专项优化整体 WAPE 下降了大约 2 个百分点从 17% 左右降到了 15% 附近。5.3 高峰期预测漂移的根本性改进高峰期预测漂移是时序预测比赛绕不开的问题。地铁早晚高峰的瞬时流量能达到平峰时段的 5~10 倍而且在很短的时间窗口内波动剧烈。我通过残差分析发现模型在高峰期的主要问题不是“预测值偏低”而是“预测值整体均值的回归”——因为训练集中平峰样本占了大多数模型天然倾向于把预测值拉向平均水平导致高峰期欠拟合、平峰期过拟合。我做了两个针对性改动。第一是把训练样本按“真实值大小”分桶在每个桶内单独计算损失并观察模型预测分布结果发现模型在低客流量站点存在明显的过估。原因是整体平均值被少数大站拉高了。第二是加入“站点客流分位数”特征也就是每个站点在全部站点中的客流规模等级配合分组建模策略让模型能区分“小站”和“大站”的不同模式。这两个改动合在一起让整体 WAPE 又下降了约 1.5 个百分点最终稳定在 12% 左右。这一步给我的启发是分布式环境下数据集可以做得很大但调优时一定要“分组看指标”而不是只看全局指标否则很容易被全局平均值掩盖掉局部的问题。5.4 分布式带来的隐藏坑标签泄漏与数据去重回到分布式主题这是整个项目里最让我后怕的一个坑。在做时间窗口特征时我需要用过去时段的真实值构造特征。在训练集上这很自然直接取上一时刻的真实值就行。但在测试集上根本不存在真实值只能用预测值回填。我一度为了省事在构造测试集特征时沿用了训练集的代码逻辑导致测试特征里用到了当天的真实值——这就是典型的标签泄漏。单机环境下数据量小这种问题很容易在抽样检查时发现。但分布式环境下数据分散在几百个分区里特征列多达 130 个某一个特征列“看起来正常其实用了未来数据”的情况非常隐蔽。我的排查方法是随机抽取几个站点的特征把同一目标时间点对应的所有特征列的取值沿时间轴可视化逐个确认所有特征都严格来自预测时刻之前。这个过程很费时间但绝对值得因为标签泄漏带来的“虚假高分”在线上根本复现不了。另外一个分布式环境下特有的问题是日期边界上的重复数据。某个 15 分钟时段的数据可能因为分区导入任务的重试机制被写入了两行完全相同的记录。在训练阶段多一条重复样本问题不大但在构造测试集特征时重复数据会让“昨日同时段值”等特征产生歧义。我在 ClickHouse 导入时改用 ReplacingMergeTree 引擎解决这个问题按(station_id, datetime)做去重确保特征计算的输入数据是唯一的。6. 复盘与落地思考分布式在城市AI计算里的真正位置6.1 分布式带来的真正价值不是“算得快”打完整场比赛回头看基于 Distributed-Platform 和 Distributed-Databases 这套架构最核心的价值其实不是把训练时间从 4 小时缩短到 40 分钟——虽然这一点也很重要。更关键的是分布式架构让迭代实验的周期被大幅压缩了。特征工程阶段我尝试了将近 60 组特征方案。在单机方案里每次全量特征重算需要几个小时一天只能做两三轮实验。分布式方案下一次全量特征重算只需要 25 分钟我就可以把一天的时间投入十几轮实验中。模型精度的最终提升很大程度上是靠“实验次数”堆出来的。这个逻辑放到真实生产环境同样成立。城市级别的 AI 计算比如交通流量预测、人群聚集预警、共享单车调度本质上都不是“一次训练出结果”的问题而是需要算法工程师和数据工程师不断迭代调优的系统工程。分布式平台的真正价值是让这种迭代以足够快的速度发生。6.2 从比赛到真实城市这套架构还缺什么比赛结束后我认真想过一个问题这套分布式架构离一个真实可用的城市 AI 计算系统还差什么首先是实时性。比赛数据是批量的但真实的城市交通系统需要接入实时客流数据预测模型需要定时滚动更新。这要求架构里补充消息队列如 Kafka和流处理引擎如 Flink把离线的批式特征计算和实时的流式特征计算结合起来做到批流一体。其次是容错与可观测性。分布式集群中节点故障、网络抖动都会影响任务执行生产环境必须有完善的任务监控、失败重试和告警机制。比赛里任务失败了重跑一遍就行生产环境里每一次失败都可能影响真实业务。第三是权限与数据安全。真实城市数据涉及大量个人信息需要在数据链路的每一层做脱敏、加密和权限控制。这一点在比赛环境里完全感受不到但在实际落地时往往是最大的工程成本。还有一个更抽象的问题是团队协作效率。算法工程师、数据工程师、运维工程师在同一个平台上的协作模式决定了这套系统能走多远。比赛里我一个人写完了从数据导入到模型训练的所有代码但在真实团队里每一步都需要清晰的接口约定和标准化的流程管理。6.3 一个最后的小技巧最后分享一个我踩了两次坑才总结出来的小技巧在分布式环境里做特征工程时一定要把“数据版本”当成一等公民来管理。我因为修改特征代码后没有清理旧缓存导致某次实验混入了旧版本的特征数据结果模型评估指标看起来很好提交后却明显下滑。后来我给每个特征宽表文件都加了生成时间戳和特征版本号训练代码强制指定版本号读取数据。这个做法虽然简单但能避免大量“莫名其妙”的实验结果波动。如果用一句话总结这段经历那就是分布式平台负责把城市数据“算得动”分布式数据库负责把城市数据“存得顺”而真正决定预测效果上限的始终是你对业务、对特征、对数据质量的理解。天池这次比赛教给我的远不止怎么调 LightGBM 参数这么简单。