news 2026/9/18 14:14:00

实时特征平台架构:美团配送的分钟级统一与Flink动态计算实践

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
实时特征平台架构:美团配送的分钟级统一与Flink动态计算实践

简介:一份美团配送实时特征平台建设实践的技术分享PDF,面向大数据实时计算开发者、平台架构师与算法工程同学。内容紧扣配送业务分钟级实时特征需求,系统梳理从平台目标、整体架构到稳定性建设、规模化的完整演进路径。包内共1个PDF文件,压缩包约56.23MB,已有195人学习下载。资料重点讲解数据输入、加工、计算、输出四层架构,通过SQL+UDF模式提升效率,以拼图式填充和上游合流应对数据乱序、实现端到端不丢不重;计算层采用基于内存计算、无状态可扩展的升级框架,支撑ETA、爆单、定价等实时特征服务。稳定性建设覆盖四层监控、隔离、双缓存、熔断限流与容灾体系;规模化部分则给出数据倾斜、高并发场景下的分片、能者多劳及本地缓存等性能优化方案。整体内容兼具架构蓝图与工程细节,可作为实时特征平台建设的系统参考。

1. 美团配送实时特征平台:从烟囱式开发到分钟级统一架构

2017 年的美团配送,履约过程正从规则驱动切换成算法驱动:调度要回答"派给哪个骑手",ETA 要预估"商家出餐多久、骑手几点到",定价要动态算"这单收多少钱",爆单要预判"哪条商圈马上拥堵"。这些算法都需要分钟级时效的实时特征,而当时每个算法团队都在自己业务系统里"烟囱式"地造特征——流程长、重复建设、稳定性还压在线上。配送数据组从 2017 年起把这件事做成了独立平台,经历了系统化、规模化、平台化三个阶段。这篇文章拆解的是三个阶段里最值得复用的设计决策:拼图式宽表如何对抗流乱序、自研 FCS 计算框架为什么长这样、50ms 响应下稳定性怎么用制度兜住、以及最终为什么又引入 Flink 做动态维度计算。

2. 系统化第一役:订单→包裹→运单的数据建模与拼图式宽表

2.1 从订单到运单:8 个核心时间点怎么抽象出来

配送履约链路长且状态散落。一次完整履约涉及用户下单、派单、骑手到店、商家出餐、骑手离店、上车、到客、用户收餐这 8 个核心时间点,中间还夹着骑手步行、驻留、骑行三种移动状态,以及室内、室外两种场景切换。如果每个算法团队各自从原始消息流里挑字段,口径一定打架。系统化阶段做的第一件事,是先把数据逻辑抽成订单、包裹、运单三层:

数据层级描述典型实体
订单用户视角的一次交易order_id,下单时间
包裹订单被拆分后的配送单元package_id,关联 order_id
运单骑手实际执行的一次履约waybill_id,关联 package_id 与 rider_id

三层关系确定后,8 个核心时间点作为标准履约模型沉淀下来,成为实时特征平台的事实标准。后续不管是 ETA、调度还是定价策略,都从这套标准时间轴里取数,而不是各写各的解析逻辑。

2.1.1 时间点抽取的具体做法

从消息流中抽取时间点,常见做法是建立一张"事件到时间点"的映射表。事件源包括订单状态机变更、骑手 App 上报的 GPS 轨迹、商家 POS 回传的出餐通知等。每一个时间点最终落到运单宽表的一列,列名统一、语义统一、更新时间由事件驱动,所有下游消费方看到的字段定义完全一致。

订单状态变更事件 -> order_time, pay_time, dispatch_time 骑手位置上报事件 -> rider_arrive_shop_time, rider_leave_shop_time 配送状态机变更 -> pick_up_time, arrive_customer_time, finish_time

这套映射在初期看起来只是"建表规范",但它的价值在三个月后才会显现:当新增一个算法模型需要"骑手到店到出餐的等待时长"这个特征,只需要在宽表模板里加一列,已经上线的计算任务全部自动兼容。

2.2 拼图式宽表:解决流乱序与端到端 Exactly-Once

实时特征平台最核心的难点不在计算,而在数据流的完整性。配送场景里,同一笔运单的各个事件在 Kafka 中的到达时间并不等价于业务发生时间,骑手在电梯里信号丢失、商家 POS 网络闪断,都会导致事件延迟甚至乱序。如果每来一个事件就直接更新宽表,下游算法很容易读到"中间态"的特征,比如订单还没有派单时间就先去算了配送时长。

方案是"拼图式"宽表。提前把运单宽表的 schema 构建成完整拼图模板,8 个时间点全部预置列,哪个事件到了就填充哪一列,其余列保持空值等待后续事件。模板里的列顺序固定,不存在"动态加列导致历史数据错位"的问题。

CREATE TABLE dwd_waybill_wide ( waybill_id BIGINT, order_id BIGINT, -- 8 个核心时间点,事件到达后按事件类型填充 order_time TIMESTAMP, -- 用户下单 pay_time TIMESTAMP, -- 支付完成 dispatch_time TIMESTAMP, -- 系统派单 rider_arrive_shop_time TIMESTAMP, -- 骑手到店 merchant_finish_time TIMESTAMP, -- 商家出餐 rider_leave_shop_time TIMESTAMP, -- 骑手离店 pick_up_time TIMESTAMP, -- 骑手取餐 arrive_customer_time TIMESTAMP, -- 骑手到客 finish_time TIMESTAMP, -- 用户收餐 -- 场景与环节标签 scene_type TINYINT, -- 0=室内, 1=室外 segment_type TINYINT -- 0=步行, 1=驻留, 2=骑行 );

这个表结构自带时间轴,查询某个运单当前处于履约的哪个阶段,只需要比较哪些列为空、哪些列已填充,不需要再用时间戳做不等值关联。

2.2.1 上游等齐、下游去重

Exactly-Once 的语义在这里拆成了两道防线。上游 Kafka 生产端按 waybill_id 分区,同一运单的所有事件进同一个分区,保证上游合流后事件顺序不丢;下游 Flink/Storm 作业用 state 记录已处理的最大事件时间,事件时间小于 state 的直接丢弃,兜住重复。

注意:拼图式宽表并不做"严格等齐",而是做"容忍缺失"。ETA 模型可以接受骑手还没到店时送达时间列为空,用特制缺失标记参与预估,而不是阻塞整条链路等待迟到事件。

3. 计算层选型:自研 FCS 与 SQL+UDF 开发模式

3.1 为什么没有直接用 Storm 或 Flink 做特征计算

2017 年时,美团内部实时计算的主流选择是 Storm 和 Flink/Spark Streaming。但配送特征计算场景有几个特殊性:特征计算逻辑密集且迭代极快,算法同学经常要改一个 UDF 就发版验证;特征量从几个涨到几十个,每个特征都要消费宽表里若干列;同时业务对计算耗时要求苛刻——每 10 分钟要处理的运单数据是千万级的。

当时行业方案的短板恰好都踩在这几个点上:

方案优势在配送特征场景的短板
Storm低延迟、原语丰富开发运维成本高,SQL 化难度大,状态管理能力弱
Flink/Spark Streaming监控运维成熟、checkpoint 稳定基于关系数据库的计算模型扩展能力偏弱,并发加不上去
基于 RPC 的自研计算内存计算、无状态、可水平扩展需要自建调度和容错,初期成本高

最终选择了第三条路:自研 FCS(Feature Compute Service)计算框架。核心设计是"基于内存计算、计算无状态、可扩展"。宽表数据加载到每个 Worker 节点的内存中,计算任务通过 MQ 分发,Worker 不保存跨任务的共享状态,扩容就是加机器。这个思路的代价是放弃了 Flink 那样的统一状态管理,换来的是"计算耗时波动极小"的确定性。

3.2 SQL+UDF:借鉴离线数仓的开发模式

实时特征平台真正提升研发效率的关键,是把离线数仓的 SQL + UDF 模式搬到实时计算里。业务团队写特征逻辑时,不直接面对流式 API,而是写一段标准 SQL,平台在编译期把 SQL 翻译成 FCS 上可执行的算子图。

public class DispatchWaitingTimeUDF extends UDF { // 计算从派单到骑手接单的等待时长 // 输入:dispatch_time, accept_time // 输出:等待秒数,若 accept_time 为空则返回 -1,代表特征缺失 public long evaluate(Timestamp dispatchTime, Timestamp acceptTime) { if (dispatchTime == null || acceptTime == null) { return -1L; } return (acceptTime.getTime() - dispatchTime.getTime()) / 1000; } }

SQL 侧的使用如下——下单到派单耗时、派单到骑手接单耗时,是调度和 ETA 模型最常用的两个基础特征:

INSERT INTO feature_waybill_dispatch SELECT waybill_id, order_id, -- 下单到派单耗时,单位为秒 DispatchWaitingTimeUDF(order_time, dispatch_time) AS order_to_dispatch_sec, -- 派单到骑手接单耗时 DispatchWaitingTimeUDF(dispatch_time, accept_time) AS dispatch_to_accept_sec, -- 是否处于配送高峰时段的小时数,直接依赖宽表字段 HOUR(order_time) AS order_hour FROM dwd_waybill_wide WHERE dispatch_time IS NOT NULL;

DispatchWaitingTimeUDF里对 null 的处理是刻意的:返回 -1 而不是 0。后续下游模型看到负数自然会把该特征视为缺失,不会误认为"派单到接单只要 0 秒"。这类细节直接决定了特征质量。

3.2.1 标准化的边界:什么特征不放进 SQL 层

SQL+UDF 模式不是万能的。对于需要迭代式计算的复杂特征(比如基于骑手轨迹聚类出的驻留点识别),仍然用 Java 原生算子写成独立计算任务,通过 MQ 写回特征存储。平台统一提供 UDF 注册、版本管理、发布审批,但引擎执行层对两种模式一视同仁。

3.3 数据倾斜治理:提前分片与"能者多劳"

实时特征计算的数据倾斜比离线更头疼,因为不能在发现倾斜后再做二次 MapReduce。FCS 的思路是:在任务调度阶段就按区域分片。配送业务天然有区域属性,一个商圈的运单量远大于另一个商圈。提前把宽表数据按 area_id 分片,每个分片对应一组 FCS Worker,分区内部再用 MQ 任务队列做"能者多劳"——哪个 Worker 处理完当前批次,就从队列里取下一个任务,避免固定分片导致的忙闲不均。

定时任务 -> FCS Worker(area=1) -> H2 索引 -> MQ 定时任务 -> FCS Worker(area=2) -> H2 索引 -> MQ 定时任务 -> FCS Worker(area=3) -> H2 索引 -> MQ

每个分片内部的计算是串行的,分片之间完全并行。宽表数据落一份到分片 Worker 的本地 H2 索引里,特征计算任务只需要从 H2 里查"本分片内已就绪的运单",然后逐批处理。

注意:提前分片不等于不做动态负载均衡。FCS 的调度器会周期性统计每个分片的积压任务量,积压超过阈值时会把该分片的剩余任务拆出一个子分片,交给空闲 Worker 处理。这个设计比单纯扩大分片粒度更精细。

4. 规模化验证:50ms 响应下的稳定性工程

4.1 服务与存储的垂直拆分:一套代码,部署隔离

特征平台从系统化进入规模化后,接入方从算法团队扩展到调度、ETA、定价、爆单四个策略中心。每个策略中心的 QPS 峰值时段差异很大:定价在午高峰、爆单在天气突变时、调度全天稳定。如果共用一个服务实例,任何一个策略的流量毛刺都会影响其他策略的响应时间。

拆分原则是"按照业务场景垂直拆分,一套代码部署隔离"。四个服务共用同一份代码仓库,但各自独立部署、独立扩缩容、独立降级策略:

服务依赖存储主要调用方
ETA 实时特征服务ETA 特征存储ETA 策略中心
调度实时特征服务调度特征存储调度策略中心
定价实时特征服务定价特征存储定价策略中心
爆单实时特征服务爆单特征存储爆单策略中心

物理资源层面,数据链路拆分得更彻底:双机房 rz 和 gh 热备,Storm 集群拆成监控、运营、履约三个独立集群;Kafka 从单机房换成 Mafka 多机房容灾;ZK、离线链路、实时链路全部隔离。这套"物理故障域隔离"避免了任何一个环节的抖动传导到交易链路。

4.2 四层监控体系与三级降级制度

稳定性不只是技术问题,更是制度问题。平台定了一个硬性指标:1 分钟响应、3 分钟定位、5 分钟恢复。支撑这个指标的是四层监控体系,从下到上每一层都有独立的告警入口:

监控层次覆盖对象典型指标
硬件监控CPU、磁盘、内存、网络磁盘使用率、网络丢包率
基础组件监控DB、缓存、MQ、ES连接池占用、消费积压
性能服务监控服务异常、超时率、QPSTP99、错误率
数据质量监控特征准确性、完备性、时效性特征缺失率、延迟分钟数

数据质量监控是最容易被忽视的一层。特征算错了比特征延迟更可怕——算法拿到一个"正常值"但实际是脏数据的特征,会直接产出错误决策。质量监控的做法是:对每个特征维护一个历史分布,实时计算的特征值偏离历史分布超过 N 个标准差时触发告警,由值班人员确认是业务波动还是数据 bug。

4.2.1 三级降级矩阵怎么定

降级分计算、服务、算法兜底三层。计算层降级:FCS 作业异常时,暂停特征计算,缓存里保留上一批计算结果;服务层降级:实时特征服务超时后,返回本地缓存的历史特征值;算法兜底:策略中心调用特征服务失败时,使用预设默认值。关键在降级矩阵的触发条件和恢复条件必须写清楚,否则降级开关就是摆设。

{ "degrade_rule": { "level": "service", "trigger": "p99_latency_ms > 50 for 60s", "action": "return_cached_value", "cache_ttl_sec": 30, "recover": "p99_latency_ms < 40 for 120s" } }

这里触发阈值是 50ms、恢复阈值是 40ms,故意留了 10ms 的迟滞区间,避免监控抖动导致降级开关反复切换。TTL 设 30 秒是为了保证兜底数据的时效性上限,超过 30 秒宁可让上游走算法默认值,也不返回太旧的特征。

4.3 查询服务性能优化:从 IO 模式到对象治理

性能要求是 50ms 响应,实际压到了 4 个 9 稳定在 40ms 以内。优化不是从 300ms 到 40ms 的线性调参,而是分了三层做减法。第一层是 IO 频次:特征查询按 waybill_id 批量分组,一次 RPC 返回 100 个运单的特征,而不是逐单查询;第二层是 IO 大小:存储层"瘦身",只保留算法实际用到的字段,宽表 30 列在查询存储里压缩成 10 列;第三层是高速 IO:在服务内存里做两级缓存(本地 Caffeine + 远端 Redis),命中率做到 80% 以上,真正的远端查询只占两成。

CPU 侧的重点是 GC。特征服务每 50ms 要响应几万次查询,对象创建速度极快,Young GC 频繁会导致 STW 波动。实操中做了三件事:一是把查询结果统一用预分配的字节数组承载,避免每次 new String;二是将时间戳统一转成 long 而不是 String 传给下游,省掉解析开销;三是把高频访问的特征对象做成不可变对象,避免并发写导致的卡顿。

5. 平台化演进:事件驱动架构与 Flink 动态维度计算

5.1 为什么动态维度不能继续用 FCS

规模和稳定性问题解决后,业务提出了更多维度的特征需求。天气特征(降雨、降雪、天气等级)、骑手轨迹 GPS 实时位置、算法实时加工的特征(预计出餐时长、预计进单量)。这些特征有一个共同点:维度是动态的,不是订单或运单的固定属性。比如"当前商圈降雨量"是区域维度的特征,"骑手当前位置 500 米内未来 10 分钟预计进单量"是由时间窗口和空间范围共同决定的动态特征。

FCS 的"提前分片 + 内存计算"模型适合运单维度的批量特征,但动态维度需要连续不断的流式计算:聚合窗口、地理围栏匹配、会话拼接。再在上面硬套分片模型,要做的改造不亚于重写框架。此时引入 Flink 作为新的计算引擎,与 FCS 并存。

5.2 第三方特征接入:事件驱动 + 采集 SDK

平台化阶段的关键架构变化是把"特征来源"从平台自产扩展到了第三方系统引入。每个第三方特征源接入时,通过采集 SDK 把数据以标准化事件格式上报到 MQ 中心,特征消费端只感知 MQ 事件,不感知下游系统的内部表结构。规则是:业务系统只要发履约事件,平台就把事件转成特征;第三方平台只要提供数据流,平台就把数据加工成特征。

-- Flink SQL 作业:从第三方天气事件流计算商圈级降雨等级特征 INSERT INTO dim_area_weather_feature SELECT zone_id, MAX(rain_level) AS max_rain_level, -- 窗口内最大降雨等级 COUNT(DISTINCT report_source) AS source_count FROM third_party_weather_stream GROUP BY TUMBLE(ts, INTERVAL '10' MINUTE), zone_id;

这个 10 分钟的滚动窗口是刻意的:天气特征不需要秒级更新,10 分钟粒度既能满足 ETA 模型对天气特征时效性的要求,又能把 Flink 作业的吞吐压力控制在合理范围。窗口太小会导致大量重复聚合,太大则会让特征滞后于天气变化。

5.2.1 引擎路由:让 Flink 和 FCS 各管一段

引入 Flink 后,平台并没有把计算层统一到单一引擎,而是做了一层引擎路由。运单维度、批量可枚举的特征走 FCS;事件驱动、窗口聚合、动态维度特征走 Flink。路由规则配置在元数据管理系统里,新特征上线时声明特征类型,路由层自动分配计算引擎,算法团队不需要关心 SQL 最终跑在哪个框架上。

存储层也做了同样的拆分:FCS 产出的特征写入原有的 Redis/ES 特征存储,Flink 产出的动态维度特征写入独立的"第三方特征存储",两个存储之间不互相读写,避免异构数据的耦合。

6. 从这份实践里可以直接抄走的四个设计决策

6.1 拼图式宽表优于"来一个事件更新一行"

这套方案经过三个阶段验证仍然成立。核心原因是它把实时数据的乱序问题从"治理"变成了"容忍"——宽表模板预先定好,事件到了就填充对应列,不到了就空着,下游模型显式处理缺失值。对比"按事件主键做 merge 更新",拼图式的优势在于列与列之间没有耦合,一个事件的迟到不会阻塞其他特征的产出。实现时可以加一个"事件水位线列"记录宽表最新事件到达时间,方便监控特征新鲜度。

6.2 计算框架要按"状态维度"选择,不是按名气

FCS 和 Flink 的并存不是架构洁癖,而是计算模型不同。FCS 适合"有限分片 + 批量特征",特征是预先可知的列集合;Flink 适合"无限流动 + 动态聚合",特征维度需要窗口或地理计算。决策标准只有一条:特征计算的 key 是静态的(运单号)还是动态的(商圈+时间+天气)。静态 key 用分片计算,动态 key 用流式引擎,不必追求"一个框架解决所有问题"。

6.3 降级矩阵要写触发条件和恢复条件

规模化的稳定性建设里,最容易失效的不是监控而是恢复。很多降级开关只定义了"什么时候降级",没有定义"什么时候恢复",导致一次故障后系统长期运行在降级模式。美团的做法是给每个降级级别配上迟滞区间——降级触发阈值为 50ms,恢复阈值为 40ms 并持续 120 秒,确保系统不会在阈值边界来回振荡。这个参数组合可以直接套用到自研的特征服务里。

6.4 对象治理是 50ms 响应绕不过去的一关

GC 调优的效果往往被低估。特征服务单机处理几万 QPS 时,最直接的耗时来源不是 CPU 计算而是 Young GC 的暂停。减少对象创建、控制对象大小、避免在热路径上做字符串拼接,这三条对任何高并发查询服务都适用。一个可实操的验证方法是:压测时打开 JVM 的 GC 日志,统计 Young GC 的间隔和耗时,如果单次 GC 时间超过 5ms,先查热路径上有没有不必要的对象分配,再考虑调整堆大小。

本文还有配套的精品资源,点击获取

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/18 14:12:45

五周 2 万 Star 的 AnyDoc,TaoToken 发 Key 给 LLM

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/18 14:12:25

PyWxDump 实战指南:微信数据库解密与微信密钥获取

PyWxDump 实战指南&#xff1a;微信数据库解密与微信密钥获取 【免费下载链接】PyWxDump 删库 项目地址: https://gitcode.com/GitHub_Trending/py/PyWxDump 一条命令&#xff0c;微信数据库密钥到手&#xff0c;MSG.db 里的聊天记录、群信息、文件索引全部能解开。PyWx…

作者头像 李华
网站建设 2026/9/18 14:11:47

OTDR与GIS融合的光纤智能监控:从长度域到地理域的故障定位实践

简介&#xff1a;一份面向光纤网络运维与智能监控方向的研究文献&#xff0c;内容聚焦基于GIS和OTDR的光纤智能监控系统设计&#xff0c;尤其针对航天发射场等关键场景的光纤线路维护需求。系统将地理信息系统的空间定位能力与OTDR实时监测能力结合&#xff0c;实现光纤故障快速…

作者头像 李华
网站建设 2026/9/18 14:09:55

组织效能分析自动化:从PPTX报告到可复算系统

简介&#xff1a;这份PPT研究报告聚焦组织效能的底层逻辑、方法论框架与案例解析&#xff0c;面向企业管理者、HR从业者及组织发展顾问&#xff0c;帮助其系统理解组织效能内涵&#xff0c;并从经营、运营、人力三个层面找到效能提升的撬动点。内容涵盖组织效能的核心要素与影响…

作者头像 李华
网站建设 2026/9/18 14:09:33

LVGL多页面切换实战:从界面编辑器到代码整合全攻略

做嵌入式GUI开发的朋友&#xff0c;肯定都有过这种经历&#xff1a;界面从设计稿到真机&#xff0c;中间隔着一条又宽又深的河。特别是搞多页面切换的时候&#xff0c;逻辑本身不复杂&#xff0c;但代码量一上来&#xff0c;页面管理的琐碎细节能把人折磨到怀疑人生。后来我接触…

作者头像 李华