简介:《美团配送实时特征平台建设实践》是一份针对实时计算与数据平台方向的技术分享PDF,面向数据架构师、实时计算工程师及算法团队,系统阐述配送场景下分钟级时效实时特征平台从0到1的建设路径,解决烟囱式开发、重复建设与稳定性风险高的核心痛点。内容覆盖平台目标与整体架构、数据流处理、计算层设计、实时特征服务、稳定性建设及规模化演进,重点介绍了订单-包裹-运单的数据建模、SQL+UDF开发模式、拼图式合流解决乱序与端到端精确一次、基于内存计算的可扩展计算框架,以及四层监控、三层降级、双机房容灾等稳定性实践。整份资源为1个PDF文件,大小约56.23MB,已有195人参与学习浏览。读者可从中借鉴美团在实时特征收敛、数据质量监控、查询性能优化和平台化架构升级方面的落地经验,对规划自研实时特征平台或优化现有实时链路具有直接参考价值。
1. 实时特征平台:美团配送把特征从"按天算"压到"分钟级"的落地样本
做实时特征的团队大多有过这种经历:算法同学提需求时张口就是"我要今天每小时的区域进单量",一听业务没错,但落到数据链路就成了黑匣子。数据从业务库同步到离线数仓要 T+1,Kafka 里的实时流又没人接,最后算法只能硬着头皮用日级特征顶上去。这份《美团配送实时特征平台建设实践》讲的就是怎么把这层窗户纸捅破——用订单、包裹、运单三层建模刻画履约全程,把特征生产收敛到一个分钟级时效、支撑调度/ETA/定价/爆单四类核心策略的平台里。适合正在做实时特征、实时数仓或算法特征平台的工程师看,尤其是那些卡在"烟囱式开发、口径对不齐、稳定性没人背锅"阶段的团队。看完能拿走的是:分片计算框架怎么防倾斜、四层监控与三层降级怎么设计、查询服务怎么压到 50ms 内。
2. 划边界定架构:订单/包裹/运单三层建模与拼图式宽表
2.1 四元关系与 8 个核心时间点:先定"刻画面"再写代码
开发实时特征平台最大的坑不是技术选型,而是口径。美团配送在系统化阶段做的第一件事,是把业务抽象成用户、商家、骑手、平台的四元关系,再把履约过程拆成用户下单、支付、派单、骑手到店、商家出餐、骑手离店、到客、送达这些环节。最终沉淀为 8 个核心时间点、2 个履约场景、3 个环节的实时刻画。这听起来像业务梳理,其实是技术决策的前提:特征维度到底按订单、运单还是包裹来建模?答案在数据的自然层级里——一个订单可能拆成多个包裹,一个包裹对应一个运单,运单才是骑手履约的最小单位。
很多团队一上来就照搬离线数仓的星型模型,把实时特征做成大宽表一把梭,结果就是字段膨胀、口径冲突、上游表结构一改全线崩。美团的思路是先把订单、包裹、运单的层级关系固化下来,所有实时特征都挂在这套模型上。我先给一份当时建模维度与粒度的参数表,便于理解后续分层设计:
| 建模对象 | 主键粒度 | 核心时间点归属 | 刻画内容 |
|---|---|---|---|
| 订单表 | order_id | 下单、支付、发单 | 用户侧履约意图 |
| 包裹表 | package_id | 调度、接单、取餐、送达 | 拆单与运单生成 |
| 运单表 | shipment_id | 到店、离店、到客、签收 | 骑手侧履约轨迹 |
| 运单扩展表 | shipment_id + 时间片 | 全程时间点聚合 | 实时计算的增量结果 |
这四张表是数据层的骨架。订单表管用户侧意图,包裹表管拆单逻辑,运单表管骑手轨迹,扩展表存实时计算出来的衍生特征。注意扩展表不是简单加字段,它是为了承载每分钟级计算结果的增量写入,避免把运单表本身撑爆,也为后面"兜底降级"留了空间。
2.2 拼图式宽表:把乱序流变成可填空的模板
实时数据流最头疼的问题是乱序。一个运单的"骑手到店"时间点可能比"用户下单"先到 Kafka,直接拼接字段会有大量空值和错位。美团的做法是"拼图式开发":先在存储层把宽表模板建好,模板里所有时间点列预先定义,上游每个环节各写各的列,谁到谁填,填完一块拼图就算完成一部分。关键在合流层——上游通过主键把所有相关流合成一条,保证数据不丢;到了下游存储层,再用唯一键约束和去重逻辑解决重复写入。
这套设计在代码层面的落地方式主要是两条:合流阶段的消息路由,以及存储阶段的 Upsert。合流要注意的是 Kafka 分区策略:必须用运单号、包裹号这类业务主键作为 partition key,否则同一个运单的到店事件和离店事件被发到不同分区,下游 join 会跨分区拉数据,时延和吞吐全崩。
-- 宽表模板定义(示意,按业务主键分区) CREATE TABLE shipment_wide ( shipment_id STRING, order_id STRING, package_id STRING, user_order_time TIMESTAMP, -- 用户下单 user_pay_time TIMESTAMP, -- 支付 dispatch_time TIMESTAMP, -- 调度发单 rider_accept_time TIMESTAMP, -- 骑手接单 rider_arrive_poi_time TIMESTAMP, -- 骑手到店 poi_finish_time TIMESTAMP, -- 商家出餐 rider_leave_poi_time TIMESTAMP, -- 骑手离店 rider_arrive_cust_time TIMESTAMP,-- 骑手到客 finish_time TIMESTAMP, -- 送达 update_time TIMESTAMP, PRIMARY KEY (shipment_id) NOT ENFORCED ) WITH ( 'connector' = 'upsert-kafka', 'topic' = 'shipment_wide', 'properties.bootstrap.servers' = 'kafka1:9092,kafka2:9092', 'key.format' = 'json', 'value.format' = 'json' );这段 DDL 说明几个关键点:一是所有时间点列在源流创建时就全部定义好,后续各环节只做按列填充,不做表结构变更;二是用 upsert-kafka 连接器,天然支持相同主键的重复消息去重;三是分区键由表主键隐式决定,保证同一运单的所有事件落到同一分区。实际生产里我一般还会加一列source_system标记事件来源,方便排查"某一列长时间没被填充"是哪条链路上游没接上。
2.3 SQL+UDF 开发模式与 DWD/DIM 分层:从烟囱式到标准化
美团把离线数仓的 SQL+UDF 模式搬到了实时链路,用 DWD、DIM、宽表索引三层来约束开发规范。DWD 层做清洗和转换,DIM 层做维表建模与合流,宽表索引服务负责把明细特征聚合成服务可读的数据结构。这么做最直接的收益是开发效率:业务团队不用再各自写一套从 Kafka 到 Redis 的链路,只需要按模板提交 SQL 和 UDF,平台负责调度和资源分配。
UDF 的使用也很有讲究。像"预计出餐时长""预计进单量"这类算法实时加工特征,不适合在 SQL 里硬写几百行 case when,而是封装成 UDF 扔进计算框架。一个典型的 UDF 要处理"事件缺失"的默认值问题——例如商家出餐时间没到,预计出餐时长特征不能返回 null,而要返回一个基于历史分位数的兜底值,否则下游 ETA 策略模型会直接报错或产出极端结果。
// Java UDF:计算运单当前环节时长,兜底历史分位数 public class StageDurationUDF extends ScalarFunction { public long eval(String stage, long eventTime, long defaultValue) { if (eventTime <= 0) { // 事件未到达,用历史P50兜底,避免下游拿到null return defaultP50(stage); } long now = System.currentTimeMillis(); if (now - eventTime < 0) { // 时钟乱序:丢弃这种脏数据,直接发到旁路日志 return -1L; } return now - eventTime; } }这个 UDF 里有两个容易被忽略的设计:事件时间兜底和乱序时间戳识别。返回 -1 不是错误,是要让下游 SQL 按异常值过滤并旁路记录;返回历史 P50 是"能者多劳"之外的另一层兜底思路——宁可给一个保守估计,也不让算法拿到空值。参数defaultValue由维表配置下发,不同区域能配置不同分位数,避免一刀切。
3. 计算层改造:无状态分片计算与"能者多劳"防倾斜
3.1 为什么不用纯 SQL 做实时特征:Storm 的边界与瓶颈
用 Storm 做实时特征计算,第一反应是用 Trident 或 Storm SQL。但实践下来,Storm SQL 化难度很高,而且基于关系数据库的计算模型扩展性差——特征计算需要大量状态查询和维表关联,关系模型很难表达"按业务分片、跨分片汇聚"这类实时需求。开发运维成本也高,团队要同时维护拓扑、消息语义和状态后端,稳定性全靠老师傅的"手感"撑着。
另一个被反复验证的教训是:实时特征计算不能完全依赖外部存储做关联。早期方案是把维表放 Redis,每个事件来都去查一次,吞吐一高 Redis 先扛不住,接着 Storm 拓扑背压,最后 Kafka 堆积。美团的解法是把计算做成"无状态 + 内存分片":所有数据按业务维度提前分片,每个分片内的计算都在本地内存完成,没有跨节点的状态访问。分片信息提前配置在 Worker 的本地文件里,运行时只做"查分片配置 → 处理本地分片数据 → 输出结果"这三件事。
3.2 基于业务 ID 提前分片:区域维度与运单维度的双分片策略
我刚才说按业务 ID 分片,这里要展开讲,因为分片 key 选错,数据倾斜问题会直接毁掉整个链路。美团配送的典型特征是按"商家、区域"为维度,按"订单、运单、包裹"为粒度。区域和商家天然有热点——核心商圈的单量是冷门区域的几十倍,如果把区域 ID 直接当分片 key,热区域所在的 Worker 会忙死,冷区域 Worker 空转。
做法是"两级拆分":先按大区做粗分片,再在粗分片内部按运单 ID 做细分片,分片数预先配置,与 Worker 数解耦。这样即使某个大区单量突增,细分片也能把负载匀开。同时,分片规则不是程序里写死的字符串拼接,而是配置在元数据管理系统里,运营同学可以按天调整分片数,适应节假日单量变化。
# 分片路由伪代码:分片key选择与Worker调度 SHARD_CONFIG = { "region": ["r1", "r2", "r3"], "shard_per_region": 8, "worker_num": 24 } def route_key(order): region = order["region_id"] shard = hash(order["shipment_id"]) % SHARD_CONFIG["shard_per_region"] # 分片名 = 区域 + 分片序号,保证同一运单进同一分片 return f"{region}_{shard}"这里shard_per_region是防倾斜的核心参数。之前线上出过一次翻车:某热门区域一个分片扛了全区域 40% 的流量,FCS Worker CPU 打满,特征产出延迟从 40s 涨到 5 分钟。后来把shard_per_region从 8 调到 16,问题立刻缓解。注意route_key里用shipment_id而不是order_id做 hash 对象,因为一个订单拆出多个包裹后,如果按订单 hash,同订单的多个包裹会被路由到同一分片,照样倾斜。
3.3 "能者多劳"模式:Work Stealing 思想在特征计算里的落地
分片只能把数据尽可能均匀铺开,但运行时各分片的计算量不可能完全一样——某个包裹轨迹特别长、事件特别多,它的分片就是比别人慢。美团的"能者多劳"模式本质是带抢占的任务队列:每个 Worker 维护一个任务队列,队列里是待处理的分片数据;Worker 处理完自己队列里所有分片后,不是空等,而是主动从其他队列偷任务过来处理。
这套机制写起来不复杂,但有一个硬前提:计算必须无状态。如果任务处理依赖上一个分片的计算结果,偷任务会导致状态错位。所以 FCS Worker 的计算输入只依赖宽表里已落地的数据,不再依赖 Worker 本地状态。任务队列用 MQ 实现,队列名按 Worker ID 区分,偷任务就是"消费别人的队列"。注意这里的数据顺序性靠"分片队列 + 分片内有序"保证,同一分片的数据只进同一个队列,跨分片之间不要求全局有序。
// Worker 消费与偷任务的核心逻辑 public void run() throws Exception { while (true) { ShardTask task = myQueue.poll(500, TimeUnit.MILLISECONDS); if (task != null) { compute(task); // 无状态计算,结果写入宽表索引 continue; } // 能者多劳:空闲时去偷其他Worker队列的任务 for (String workerId : peerWorkerIds) { ShardTask stolen = steal(workerId); if (stolen != null) { compute(stolen); metrics.increment("stolen_count"); break; } } } }这段代码是"能者多劳"的最小实现。myQueue.poll(500ms)是给偷任务留出机会窗口;steal要走 RPC 调用其他 Worker 的内存队列,所以每个 Worker 都要暴露一个轻量的取任务接口。实际生产我给stolen_count加了监控,如果它长期为 0,说明分片配置过粗,需要调大shard_per_region。
3.4 定时任务与 H2 索引:宽表数据的最终落点
FCS Worker 计算完的结果直接写宽表索引表,索引表用 H2 内存数据库承载。为什么不直接用 Redis?因为特征计算产出的数据是"宽表行"结构,带几十个字段,Redis 的 hash 存储需要把每个字段拆成 key,序列化和反序列化开销反而更大。H2 作为嵌入式内存库,能直接用 SQL 做条件查询,还能批量更新一行中的多个字段。
索引表的作用是让下游查询服务能以"运单 ID + 时间片"的方式快速拿到特征。写入路径是"FCS 计算结果 → MQ → 索引表写入服务",这里 MQ 的解耦很关键:实时计算不再直接面向查询流量,计算的高吞吐和查询的低延迟互不拖累。
4. 稳定性与数据质量避坑:四层监控、三层降级与容量冗余
4.1 现象到根因:稳定性建设不是"加监控"而是"定制度"
先讲一个自己踩过的坑,后面大家做方案评审时能省很多事。
现象:某个实时特征突然大面积为空,算法模型线上推理结果异常,ETA 预估偏了 20% 以上。排查时监控面板一片绿色,根本定位不到是哪层出了问题。
原因:当时只做了服务层监控(QPS、响应时间),没覆盖数据质量维度。实际上 Kafka 集群内部发生 partition leader 切换,个别分区的数据延迟从秒级涨到分钟级,但服务层的 QPS 和延迟指标完全正常,因为查询服务还在正常返回缓存里的旧特征——黑匣子问题就这么来的。
解决:把监控拆成四层——硬件层(CPU、网络、磁盘、内存)、基础组件层(缓存、MQ、ES)、性能服务层(QPS、超时率、异常率)、数据质量层(准确性、完备性、延迟、容量)。每一层单独告警,数据质量层的"特征覆盖率"指标优先于性能指标。从那以后我每次上线实时特征,都会强制走一遍"索引修复"演练,确认特征数据能从源头拉回到最新水位。
再说多机房。美团用双机房(rz、gh)热备,关键链路三集群(监控、运营、履约)垂直拆分。这里有个容易踩的设计坑:多机房部署不是把同一套服务复制到两个机房就完事,而是数据写入和流量路由要按机房做切分,故障时整体切换。比如履约集群只处理履约链路的写入,流量打到另外机房时,RPC 调用要穿透到正确机房的数据分片,不能盲目做全量同步。
最后是容量规划。1.5 倍容量、定期压测是硬指标。GH 机房断电那个经典事故,就是靠双机房热备 + 容量冗余扛过去的。我见过不少团队容量规划只做"当前流量的 1.2 倍",结果双十一大促流量翻倍直接击穿,这种属于拿单量赌稳定性,不值得学。
4.2 三层降级体系:计算降级、服务降级、算法兜底
稳定性建设的核心不是防故障,而是防故障时的雪崩。美团的容灾体系里最值得抄的是三层降级设计。第一层是计算降级:实时计算链路出现问题时,自动切到离线特征兜底,特征延迟从分钟级变成小时级,但至少数据是完整的。第二层是服务降级:查询服务本身的熔断和限流,防止雪崩。第三层是算法兜底:把特征缺失时的默认行为前置到特征服务内部,不让脏数据流到上层算法。
# 降级模块配置(示意) feature_fallback: enable: true rules: - feature: "rider_pickup_duration" strategy: "P50" fallback_value: 480s when: "stream_lag > 120s || coverage_rate < 0.95" - feature: "region_eta" strategy: "offline_snapshot" fallback_value: "${offline_eta_table}" when: "compute_available < 0.8" - feature: "order_push_time" strategy: "ignore" when: "source_event_missing"这段配置说明了降级的三个策略层。stream_lag > 120s触发 P50 兜底,是最轻量的方式;coverage_rate < 0.95是数据质量监控里常被忽视的指标,覆盖率掉到 95% 以下就说明有分片在丢数据;ignore策略用于"事件确实没发生"的场景,比如还没到推送时间点的订单,降级成空值比兜底假值更安全。注意降级是分特征控制的,不是全平台统一开关。有的特征对实时性极其敏感,比如调度派单的骑手距离,这种就不能按秒级降级,要走双链路对比的旁路方案。
4.3 从 Kafka 集群故障看兜底的有效性
2018 年那个 Kafka 集群故障案例里,特征兜底避免了线上事故。当时的情况是某个 Kafka 集群出现长时间不可用,实时计算链路全部停摆。如果没有兜底,调度策略直接拿不到骑手的实时位置特征,整个派单逻辑会退化到最原始的时间片轮询。兜底策略把骑手位置特征切换到历史轨迹外推,虽然精度下降,但调度链路没断。事后复盘有一个血泪经验:兜底不是只写 default 值,还要把"兜底命中率"作为一个核心监控指标。如果兜底命中率长期高于 5%,说明实时链路本身就存在慢性问题,不能等着它变成事故才处理。
4.4 数据质量全链路监控:过程质量比结果质量先暴露问题
数据质量的监控框架分为三块。流计算时效性覆盖 FCS 的延迟、完备性、准确性;实时特征服务覆盖响应时间、可用性、容量;特征结果准确性依赖离线比对任务定时抽查。这里的关键动作是把"数据质量"做成每个特征的评价指标,而不是只看链路整体。比如每隔 5 分钟抽样一批运单,拿实时平台产出的特征值和离线清洗后的口径做 diff,偏差超过阈值就告警。这套机制能发现上游业务表结构变更、埋点日志格式错误这类"数据源头崩了但服务还活着"的问题。
4.5 数据修复时间窗:实时索引修复与离线数据修复的双轨策略
实时特征平台跑久了必然遇到这么个事:某个字段因为上游 bug 算错了,等发现时已经写进宽表索引,下游算法已经消费了好几轮。美团的解法是"实时索引修复 + 离线数据修复"双轨并行。实时索引修复是直接对索引表里的特定记录做 Overwrite,修正值立即生效;离线数据修复是重跑离线数仓任务,把历史特征表里的错误数据洗掉,保证后续离线回测和模型校验时不会用到脏数据。
这个双轨机制有几个注意点:一是修复脚本要带 commit_id 和操作人,不能让人人都能改线上特征数据;二是修复完成后要强制触发一次全链路数据质量检查,确认从 DWD 到宽表索引到特征服务全链路口径一致;三是修复不能只改数据,还要找出产生脏数据的源头 SQL 或 UDF,否则同一个坑会反复踩。
5. 查询服务性能优化:50ms 4个9的IO与GC实战
5.1 IO 频次优化:批量分组读取的收益
实时特征服务的性能目标是 50ms 内返回、4 个 9 可用性、支撑 60w+ QPS。这个量级下,单次请求的 IO 开销就是生死线。最早版本是一个特征一个查询,一次算法请求要特征服务查 10 次 Redis 或 H2,串行做 10 个 round trip,光网络耗时就去掉 40ms。优化是 IO 频次这个维度上的两件事:批量查询、分组查询。
批量查询是让特征服务把一次请求里的所有特征 ID 攒成一个 batch,到存储层做 mget 或批量 SQL。分组查询是按特征维度做拆分的另一个策略:比如骑手实时特征一组,商家特征一组,两组分别从不同的存储集群读取。这利用了数据特征的访问局部性,但注意会增加一次请求的深度——换来的是单次 IO 的体量变小,整体链路更容易控制超时。
-- 按运单批量拉取特征:拼IN查询,避免循环单查 SELECT shipment_id, feature_key, feature_value FROM shipment_feature_index WHERE shipment_id IN ( 'S202501010001', 'S202501010002', 'S202501010003' ) AND feature_key IN ('rider_poi_dist', 'poi_finish_time', 'eta_prediction') AND ts >= '2025-01-01 00:00:00' AND ts < '2025-01-01 00:05:00'SQL 里有两个参数值得关注。ts的时间窗口是"时间片"机制——特征按分钟生成版本,查询时只取目标分钟片的数据而不是最新一条,这样能规避写入乱序导致的读到半条状态;feature_key的列表要控制在 20 个以内,太多会让单条 SQL 的执行计划变得很重,反而拖慢性能。批量查询的代价是单次响应的大小增大,所以要对返回字段做瘦身——只取算法需要的字段,不取整个宽表行。线上实践里,把这个查询从"按 key 循环 get"改成拼 IN 后,TP99 直接降了 30%。
5.2 双缓存与本地缓存:把热数据留在进程内
缓存是另一个大头。美团的做法是"双缓存":本地缓存 + 分布式缓存两级。本地缓存用 Caffeine 或 Guava Cache,存的是最近几分钟内被高频访问的特征;分布式缓存用 Redis 或 H2,存全量特征。查询顺序是先查本地,命中直接返回;不命中再查远程。
双缓存有个容易踩坑的地方:一致性。实时特征本身是分钟级更新的,所以本地缓存 TTL 设置不能太长,一般 30~60s。别为了命中率把 TTL 调到 10 分钟——特征已经更新了,算法还拿着旧值,调度效果打折。另一个参数是每条特征的大小,本地缓存是堆内存,对象越大 GC 压力越大,所以本地缓存里只放 JSON 序列化后的字符串,不放对象。
# 本地缓存配置 cache: type: caffeine maximum_size: 100000 # 最多缓存10万个key expire_after_write: 45s # 45秒过期,配合分钟级特征更新时间片 record_stats: true remote_cache: type: redis batch_size: 20 timeout_ms: 15expire_after_write: 45s这个值来自"特征每分钟更新 + 容忍最多 15s 延迟"的折中。如果业务上对延迟更敏感,可以压到 30s,但本地缓存命中率会明显下降;反过来拉到 60s,命中率上来了,但算法吃到的特征可能滞后一个完整周期。batch_size: 20的意思是远程缓存读也走批量管道,这一步在高峰期能省大量 RTT。
5.3 减少对象创建:GC 停顿从 200ms 压到 20ms
实时特征服务是高并发低延迟场景,JVM GC 是大敌。早期 TP99 不达标,排查下来是 Full GC 频繁,单次停顿 200ms+。根因是两个:一是每个请求都在创建大量中间对象(尤其是字符串拼接和 SimpleDateFormat),二是缓存里存的对象过大。后面做了一轮优化,核心三板斧:用 StringBuilder 代替字符串拼接、用 ThreadLocal 复用 DateFormat、控制对象大小。
还有一个容易忽略的点:批量查询返回的结果集对象。如果一次查询返回 20 个特征,每个特征又是一个包含多个字段的对象,那单次请求就产生 100 个对象。优化办法是返回扁平结构——用Feature[]数组代替List<Feature>,用原始类型long代替Long。这些微优化每一项的收益不大,但叠在一起能把 Young GC 频率降一个量级。
// 扁平化特征读取:避免在循环里创建大量包装对象 Feature[] features = new Feature[keys.length]; for (int i = 0; i < keys.length; i++) { long value = indexReader.read(keys[i]); // 直接返回原始类型 features[i] = new Feature(keys[i], value); // 复用固定数组,不扩容 }这段代码里的new Feature(keys[i], value)是不可避免的对象创建,但数组是预分配的,不会触发 ArrayList 的扩容复制。indexReader.read()拿到的是原始类型 long,替代了装箱后的 Long 对象,减少了 GC 根扫描的负担。实际压测里,这个循环从 List 改成数组后,单请求对象创建数下降了 30%。GC 优化没有银弹,就是把每一个角落的浪费都抠出来。
5.4 容量规划与压测:1.5 倍不是拍脑袋
容量规划维度上,美团的标准是 1.5 倍容量冗余 + 定期压测。这里有三个维度要压:流计算集群的吞吐、特征存储的 QPS、查询服务的响应时间。压测不是随便拿压测工具打满就行,要按业务高峰期的特征重放,尤其是大促时段的流量模型——日常流量曲线和峰值流量的分布差异极大,拿平均值做压测会因为分片热点问题在高并发下重新暴露而失真。
压测应该产出一张"容量水位表",标出每个核心服务在当前流量下的 CPU、内存、GC、延迟四项指标。当 CPU 超过 60% 或 TP99 超过 40ms 时就要扩容,不能等到 80% 才动手,因为实时链路的故障是串联的,一个环节满了,整个链路都会垮。1.5 倍冗余的意义在于:即使一个机房断电或一个集群故障,剩余容量仍然能扛住全部流量。
6. 平台化收尾:从特征收口到事件驱动的开放式架构
6.1 垂直拆分:一套代码、多套部署的隔离落地
平台化阶段最实操的经验是垂直拆分。调度、ETA、定价、爆单四个业务团队各需要一套特征服务,但不能各搞一套代码,否则又退回烟囱式开发。最终方案是"一套代码 + 多套部署":通过配置区分服务名、存储集群、特征分组。这样既做到了资源隔离——某个业务的特征服务故障不影响其他业务,又保留了一套代码统一维护的研发效率。
拆分的边界不是按团队划的,而是按业务场景划的。四个场景的特征集有重叠但差异更大:调度关注骑手实时位置和负载,ETA 关注时间点预估和出餐时长,定价关注供需比和区域压力,爆单关注单量突增。每个场景独立部署后,容量规划也清晰了:调度场景流量大,给它 50% 资源;爆单场景只在高峰时段有压力,设置弹性扩缩容。
6.2 事件驱动与 Flink 上收:把第三方便特征纳入计算层
平台化的第二个动作是开放。原来的架构只支持平台内部自产的特征,但算法对更多粒度的特征有需求:天气(降雨、降雪、天气等级)、骑手轨迹(GPS)、算法实时加工的特征(预计出餐时长、预计进单量)。这些特征来自不同团队甚至第三方,没法都强制收敛到平台内部。
解法是事件驱动:平台开放履约事件,第三方系统通过采集 SDK 上报特征数据到 MQ;同时向上屏蔽计算引擎,引入 Flink 处理动态维度计算——比如天气特征需要按地理区域与时间窗口做 join,这种动态维度在静态分片的 FCS 里做不了。这里要说明的是 Flink 的定位:它不是替代 FCS,而是补充。FCS 继续负责高吞吐、固定分片的基础特征,Flink 负责维度灵活、窗口多变的计算特征。引擎路由层根据特征注册表决定新的特征走哪条计算链路。
// 采集SDK埋点:第三方特征上报核心代码(示意) FeatureCollector collector = FeatureCollector.getInstance(); collector.setAppName("weather_service"); collector.setEvent("weather_level_change"); JSONObject payload = new JSONObject(); payload.put("region_id", "r1"); payload.put("weather_level", 3); // 雨雪等级 payload.put("rainfall", 12.5); // 降雨量(mm/h) payload.put("event_time", System.currentTimeMillis()); // 异步批量发送,带本地缓存兜底,避免业务线程阻塞 collector.reportAsync(payload, 5000);这段 SDK 代码有两个细节。reportAsync(payload, 5000)的 5000 是本地缓冲队列长度,第三方系统在极端情况下发送速度超过平台接收速度时,SDK 会在本地阻塞而不是直接丢弃消息,给平台端留出处理时间;setAppName和setEvent是上报元数据,平台端根据这两个字段做特征分组和鉴权,防止第三方系统越权上报其他领域的特征。
6.3 验证方法:从"上线即无 S 级事故"到主动容量规划
平台化阶段做完后,我习惯用三件事验证整体状态。第一件是"特征收口率":线上还有多少特征没走平台?这个指标很直接,如果收口率低于 95%,说明还有业务团队在偷偷自建链路,平台的价值就会被打折扣。第二件是"分钟级时效达标率":每分钟产出的特征里,多少比例在 40s 内完成计算?这个指标比集群整体吞吐更敏感,能暴露分片热点和队列堆积。第三件是"故障恢复时长":从监控告警到定位根因到恢复服务,能否稳定在 3 分钟内完成?这个指标靠的是制度——值班制度、巡检制度、报警治理、Case Study 总结,技术只是底子。
最后一个技巧,也是我在做容量规划时最常用的一条:不要只关注峰值 QPS,要看"每分钟特征生产量"这个指标。美团高峰期每分钟生产 1000w+ 特征、计算耗时控制在 40s 以内,这意味着如果特征生产量涨了而计算耗时没变,说明系统还有余量;如果生产量涨了、耗时也跟着涨了,说明资源要到瓶颈了,该提前扩容。我习惯每季度做一次容量水位复盘,把那段时间的特征生产峰值、计算耗时、查询 TP99 画在一张表里,数据一多,系统性的资源瓶颈自然就浮现出来了。从那以后我每次接新的特征需求,都会先问一句"这个特征走哪条链路、兜底策略是什么、覆盖率指标挂在哪个面板上",三个问题答清楚才动手开发,希望帮到你。
本文还有配套的精品资源,点击获取