恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
美团配送实时特征平台:从事件接入到在线离线一致性
首页
资讯中心
/
美团配送实时特征平台:从事件接入到在线离线一致性
美团配送实时特征平台:从事件接入到在线离线一致性
发布时间:2026/9/18 15:11:55
简介来自美团配送团队的实时特征平台建设实践资料面向大数据平台、实时数仓及算法工程方向的工程师与架构师用于解决分钟级实时特征建设与稳定性治理难题。资源包仅含 1 个 PDF 文档大小 56.23MB支持离线研读。目前已有 195 人学习下载。内容围绕平台建设从系统化到平台化的演进展开覆盖平台目标、整体架构、数据流清洗与合流、SQLUDF 计算层、实时特征服务等核心模块并结合调度、ETA、定价、爆单等策略场景详解了防数据倾斜的分片任务、三层降级兜底、双机房容灾、全链路监控以及 TP4 个 9 稳定在 40ms 内的查询性能优化做法。对希望构建或升级实时特征平台的团队而言其中沉淀的架构设计与稳定性治理经验具有较强的工程参考价值。1. 美团配送实时特征平台的起点实时特征不只图快日常配送调度中骑手侧每看到一个推荐订单服务端往往已经用近 5 到 10 分钟的接单率、取消率、骑手轨迹密度等特征完成了一次在线推理。美团配送实时特征平台要解决的核心问题就是为这类模型准备秒级可见、可回放、口径稳定的事件特征订单从创建到被接单花了多久、骑手是否正在绕路、取消是否在某个商区集中这些状态变化不能等离线任务跑完再进入模型。这篇文章适合负责服务端链路或大数据研发的同学阅读目标是梳理从事件接入、实时计算、特征存储到时间一致性兜底的一整条落地路线。先定一个判断这类平台通常围绕 Kafka、Flink 和 KV 存储来组织但真正拉开差距的是如何统一在线与离线特征口径以及如何处理配送 App 在弱网下天然存在的迟到事件。2. 配送事件接入与特征定义把订单状态流转成可计算特征实时特征平台的第一层不是计算引擎而是事件模型。配送场景里有两类事件天然适合驱动特征计算一类是订单状态变化另一类是骑手轨迹上报。它们都存在业务表和日志里但不能直接拿原始表去算需要先转换成带语义的事件流再定义可配置的特征。2.1 配送事件模型从订单状态机提取特征源动词一张配送订单表本质上是一行被不断 UPDATE 的记录直接监听表变更很难看出“接单后多久取消”这类迁移过程。常见做法是解析 binlog 后将订单状态迁移转换成语义事件例如ORDER_CREATED、ORDER_ACCEPTED、PICKUP_COMPLETED、DELIVERED、ORDER_CANCELLED。每条事件 payload 中携带pre_status和cur_status这样特征配置时可以直接按event_type过滤不需要在 SQL 里反复解析业务状态。骑手轨迹流则相对简单由骑手 App 按秒级或百毫秒级上报经纬度、速度和上报时间。这类事件不涉及状态机但数据量远大于订单事件是平台里最容易产生成本的一类输入。两股流在进入特征平台前需要统一两个字段rider_id作为主键event_time作为事件时间。event_time必须取业务实际发生时间而不是服务端接收时间否则后续时间窗口全部会对不齐。事件源代表事件示例特征常用窗口订单状态流ORDER_CREATED、ORDER_ACCEPTED骑手接单率、平均接单时长10 分钟滑动轨迹流GPS_UPLOAD最近 5 分钟平均速度、累计位移5 分钟滚动取消事件流ORDER_CANCELLED商圈取消率、取消订单骑手状态30 分钟滑动商圈围栏事件RIDER_ENTER_POI是否处于配送热点区域瞬时状态上表并不是把原始事件直接入库而是由平台解析出可复用的事件源。事件类型定义得越细特征配置阶段就越少出现“同一个字段在不同任务里含义不同”的问题。2.2 特征配置语言一份 YAML 把窗口、过滤、聚合都描述清楚实时特征一旦多了之后最常见的失控方式是在多个 Flink 任务里各写各的 SQL同名特征在不同任务里窗口大小不一致或者过滤条件差一个event_type。我一般会先用一层特征配置把口经固定下来再让平台把配置编译成 Flink SQL。配送场景下可以这样描述一个特征feature_id: rider_order_accept_rate_10m owner: dispatch-core source: delivery_order_event key_by: rider_id window: type: sliding size: 10m slide: 1m allowed_lateness: 5s metric: name: accept_rate value: accept_cnt / order_cnt agg: ratio filter: event_type in (ORDER_ACCEPTED, ORDER_CREATED)这份配置的可读性不需要解释关键在几个字段的约束owner标识负责人出现口径争议时知道找谁window.slide决定特征刷新粒度配送场景用 1 分钟作为 slide 比用 5 更平滑allowed_lateness会直接传给 Flink 窗口表示允许迟到多久。平台构建时检查配置是否引用已注册的事件源如果发现新增事件类型未在事件字典中构建直接失败避免线上任务静默读取不存在的字段。配置语言本身不用很复杂。它的价值是成为在线任务和离线任务的唯一元数据来源在线任务读它生成 Flink SQL离线回填任务读它生成 Hive SQL。这样在线特征和离线特征天然同源回放和验证才有基础。2.3 用 Flink SQL 把配置翻译成第一个实时特征任务拿到上面rider_order_accept_rate_10m后平台会生成类似下面的 SQL。这里使用 Flink SQL 的 HOP 窗口表达 10 分钟滑动、1 分钟步进。CREATE TABLE delivery_order_event ( order_id STRING, rider_id STRING, event_type STRING, event_time TIMESTAMP_LTZ(3), WATERMARK FOR event_time AS event_time - INTERVAL 5 SECOND ) WITH ( connector kafka, topic delivery_order_event, properties.bootstrap.servers kafka:9092, properties.group.id feature_platform, format json ); CREATE VIEW rider_accept_cnt AS SELECT rider_id, HOP_START(event_time, INTERVAL 1 MINUTE, INTERVAL 10 MINUTE) AS win_start, COUNT(*) AS accept_cnt FROM delivery_order_event WHERE event_type ORDER_ACCEPTED GROUP BY rider_id, HOP(event_time, INTERVAL 1 MINUTE, INTERVAL 10 MINUTE); CREATE VIEW rider_order_cnt AS SELECT rider_id, HOP_START(event_time, INTERVAL 1 MINUTE, INTERVAL 10 MINUTE) AS win_start, COUNT(*) AS order_cnt FROM delivery_order_event WHERE event_type IN (ORDER_ACCEPTED, ORDER_CREATED) GROUP BY rider_id, HOP(event_time, INTERVAL 1 MINUTE, INTERVAL 10 MINUTE);这里需要重点说明几个参数的取舍。HOP的第一个参数是事件时间第二个参数是 slide第三个参数是 size表示每 1 分钟向前滚动一次每次计算过去 10 分钟内的聚合值。WATERMARK设置了 5 秒意思是允许最多晚到 5 秒的事件参与窗口超过这个时间才被判定迟到。配送业务中服务端接收时间和骑手实际发生时间的偏差通常在几秒以内因此 5 秒是一个偏安全的默认值如果特征是“是否取消”这类强时效信号可以收紧到 2 到 3 秒。两个 VIEW 对应一个分子的一个分母最终accept_cnt / order_cnt可以在外层 join 算出也可以直接在 Flink 任务里再套一层查询。平台会尽量将这类多窗口计算合并进同一个作业但状态后端仍然按 rider_id 分组因此需要确保两个视图使用同一个 Kafka group.id避免重复消费导致两边时间戳不一致。2.4 特征注册与血缘平台层和计算层要剥离开如果每增加一个特征都上线一个独立 Flink 作业体量变大后运维成本会压垮整个平台。特征配置入库后平台层保存的是特征元数据、负责人、上游事件源和窗口参数执行引擎只负责运行由元数据生成的 SQL。每次变更特征时先走配置变更平台自动生成 diff确认只影响新增维度或窗口后热更新配置不重启作业。血缘关系也随之建立任何一个特征变更都能反查它影响了哪些在线推理链路。离线回填任务复读同一份配置这是后面做一致性验证的关键前提也是一个平台叫“平台”而不是“一套 Flink 任务”的底气。3. 实时特征存储与读取路径HBase、Redis 双写保住低延迟特征计算完成之后不能让业务推理直接查 Flink 状态必须把结果物化到存储层。配送场景对特征读取的延迟要求很苛刻模型推理时一次要取多个特征存储访问耗时通常不能超过个位数毫秒。这里就引出实时特征平台区别于普通流计算任务的另一条主线存储选型和双写设计。3.1 特征存储选型哈希表、时序存储和可扫描存储的分工不同特征对延迟、历史回溯和扫描方式的要求完全不同不应该只用一种存储。结合配送场景可以把特征分成三类高频长尾的聚合特征、按时间序列组织的地点特征、以及超高频访问的实时标签位。常见选型如下特征类型典型特征存储选型为什么这样选聚合特征骑手 10 分钟接单率、平均配送时长Redis Hash读延迟最低适合模型在线访问序列特征最近 10 个 GPS 轨迹点、最近 30 分钟位移HBase按 rider_id 和时间戳扫描实时标签是否在商圈、是否有在途订单Redis Bitmap / 内存超高频读取要求极低延迟历史回放数据某天所有窗口的特征快照HBase 或对象存储用于离线回放和 validation配送模型的一次请求往往要同时读取多个特征因此我倾向于把同一骑手的所有聚合特征放在同一个 Redis Hash 中而不是拆成多个字符串 key。这样一次网络往返能取回全部所需字段减少模型服务的连接和时间消耗。HBase 则用来保存轨迹序列和特征多版本历史因为 Redis 不适合做范围扫描和批量回放。3.2 Redis 的 Hash 设计与写入幂等Redis 侧的结构一般写成feat:v1:rider:{rider_id}:10mfield 为 feature_idvalue 为聚合结果的字符串序列化。同一个骑手会在一个 key 下保留多个特征配合 expire 控制有效期。写入时Flink 任务输出的是已经聚合好的最终值Redis 端不需要做累加。但如果 Kafka 因为消费重启发生重放同一条事件可能被写入两次因此写入端必须保证幂等。把写入封装成 Lua 脚本是配送链路里比较常用的做法local current redis.call(hget, KEYS[1], ARGV[1]) if not current then redis.call(hset, KEYS[1], ARGV[1], ARGV[2]) redis.call(expire, KEYS[1], ARGV[3]) return 1 end redis.call(hset, KEYS[1], ARGV[1], ARGV[2]) return 0脚本参数依次为 key、feature_id、value、过期时间。逻辑很简单字段不存在时设置值并设置过期时间存在时直接覆盖。不要在 Lua 里做累加操作因为 Kafka 重投时同一条事件会再进来一次如果这里对计数做INCR特征值就会被重复计算。幂等应该由 Flink 作业在写入前通过event_id去重Lua 只负责最后一步的原子覆盖。3.3 HBase 表结构与 RowKey 设计按骑手聚合还是按特征聚合HBase 主要承担两类职责轨迹序列存储和特征版本回溯。轨迹序列的 RowKey 建议以rider_id开头因为一次查询总是一个骑手一个时间段的多条记录用 scan 可以快速取回。特征回放表则要考虑同一骑手在多个窗口下的多个版本RowKey 可以设计为v1_{rider_id}_{feature_id}_{hour_bucket}自然按周一小时分桶减少写入热点。建表时可以这样创建基础表create feature_store, {NAME f, VERSIONS 3, COMPRESSION SNAPPY}VERSIONS 3是特意设置的。HBase 默认只保留最新一个版本而实时特征如果发生迟到重算旧值依然很有价值通过对比相邻版本之间的变化量可以判断特征是否发生异常跳变也可以支撑特征回放时的历史比对。压缩选择 Snappy 是为了在读写速度和压缩率之间取一个平衡配送场景的高峰写入量很大不压缩很容易把 RegionServer 的 IO 打满。写入 HBase 时应手动指定时间戳把特征的事件时间而不是处理时间作为版本。这样在回放和修复时读到的是按业务时间排序的版本序列而不是按写完时排序对后续 diff 非常关键。3.4 在线读取的降级链路Redis 优先HBase 兜底模型服务读取特征时我会让读取路径按照 Redis、HBase、默认值三级逐步降级。Redis 未命中不一定代表没有特征可能是 key 过期或写入延迟所以应继续尝试从 HBase 取最近版本。取不到时才返回默认值并同时上报一条 feature_miss 指标用于观测平台覆盖率。需要注意的是业务读取端要区分“特征不存在”和“特征过期”新建骑手可能确实没有历史特征这时直接使用默认值而“过期”意味着计算链路已经有一段时间没写入需要触发异步刷新并把结果回填。落库和读链路的延迟底线可以量化为Redis p99 小于 2 毫秒HBase 兜底 p99 小于 8 毫秒特征平台整体对模型的影响控制在 10 毫秒以内。超过这个预算时优先检查的是读取端串行查询次数而不是存储本身。4. 配送弱网下的事件时间一致性水印、迟到重算与特征兜底实时特征平台最隐蔽的问题发生在弱网环境。骑手在电梯里、地下停车场或偏远路段App 上报的轨迹会积压网络恢复后这些点集中到达服务端。如果平台按处理时间计算窗口这些迟到点就会被分到错误的窗口在线特征和离线回放结果自然对不上。这一章专门讨论如何配置事件时间参数以及迟到事件到达后如何做特征补偿。4.1 迟到事件不是偶发而是配送高峰期的常态配送场景的事件迟到有两个来源一个是骑手 App 上报端的本地缓存另一个是服务端在高峰期消费积压。前者可能造成数十秒到几分钟的延迟后者相对可控。这两类迟到在午高峰、恶劣天气单量突增时会更明显。如果特征平台只在建表时把 watermark 写成固定 5 秒高峰期迟到的轨迹就会被大量丢弃模型拿到的特征在峰值时段会趋于“空白”。因此在设计阶段必须把迟到容忍度作为特征配置的一部分而不是每个任务各自拍脑袋。4.2 Flink 的 watermark 与 allowedLateness 参数组合Flink DataStream API 下一个兼容配送实际场景的窗口计算通常长这样DataStreamGpsPoint gpsStream ...; OutputTagGpsPoint lateTag new OutputTagGpsPoint(late-gps) {}; SingleOutputStreamOperatorAverageSpeed result gpsStream .assignTimestampsAndWatermarks( WatermarkStrategy.GpsPointforBoundedOutOfOrderness(Duration.ofSeconds(20)) .withTimestampAssigner((point, ts) - point.getEventTime())) .keyBy(GpsPoint::getRiderId) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .allowedLateness(Time.seconds(10)) .sideOutputLateData(lateTag) .aggregate(new AverageSpeedAggregate());forBoundedOutOfOrderness(20)表示 watermark 等于当前已观测到的最大事件时间减去 20 秒。换句话说事件时间小于这个阈值的窗口会被触发计算。allowedLateness(10)允许窗口在触发之后继续保留状态 10 秒这期间如果再有迟到事件窗口会重新计算并输出新的结果。两者的叠加含义是迟到 30 秒以内的事件都有机会修正特征超过 30 秒则进入late-gps侧输出流。参数典型配置说明watermark 延迟10 到 30 秒根据端到端延迟 p95 估算太小丢弃多太大增加窗口重算量allowedLateness等于或略小于 watermark 延迟窗口关闭后允许重算的时间长度状态 TTL窗口长度 迟到容忍防止状态保存过久控制后端内存侧输出流watermark allowedLateness 之外的极晚事件不直接丢弃交给补偿任务处理参数设置上我一般会让allowedLateness不超过 watermark 延迟否则窗口状态会保存过久而且两次重算之间的跳变幅度也会更大。侧输出流不是直接丢弃而是把极晚事件写到一个单独的主题由补偿任务读取。4.3 迟到重算后的特征补偿只依赖 Flink 重算还不够窗口重算后 Flink 会写出一个新的特征值。但这里有一个容易被忽略的问题模型在 10:02 读取到的平均取餐时长是 3 分钟10:03 一个 09:58 的迟到事件触发重算后变成了 4 分钟如果存储直接覆盖新值模型会看到一个明显的跳变。这种跳变对在线模型的影响往往大于特征本身不准确的影响。所以存储层不能简单覆盖旧值。我看到的配送场景实践是采用双缓冲当前在线特征值保持不动重算后的新值先写入新区域并带上feature_ts。模型读取时判断两个时间戳之间的差值如果差值大于阈值就对新旧值做平滑加权权重随事件时间差逐渐偏向新值。对于取消率这类高频变化的特征甚至可以忽略迟到超过一个完整窗口的旧事件只接受当前可解释范围内的重算。而那些必须精确对齐离线累积口径的特征比如当日骑手单量则必须做补偿否则离线回放任务永远对不上。4.4 读取端的特征新鲜度约束过期比缺失更危险即使计算端配置得再完善平台也可能会出现几分钟的中断或消费堆积。读取端这时候如果照常返回旧值模型不会报错但决策依据已经失真。更稳妥的做法是在 SDK 层给每个特征声明有效时间范围超出时返回默认值并触发异步刷新。feature_freshness: rider_order_accept_rate_10m: ttl: 60s fallback: 0.5 rider_location_latest: ttl: 20s fallback: null上面 YAML 的含义是读取rider_order_accept_rate_10m时如果该特征的时间戳已经超过 60 秒不再返回旧值而是用 0.5 作为默认值rider_location_latest则最多接受 20 秒的延迟对骑手实时位置来说这个值已经足够严格。这样做的代价是特征缺失率会增加但换来的是模型不会静默使用过期特征。缺特征和特征过期应该分开监控缺是正常稀疏过期则意味着链路故障。5. 实时特征平台的回归验证旁路对比收敛上线差异上线一个新的实时特征最怕的不是某个值算错而是线上在线读取的特征和离线回放对不上。对不上的原因通常不是计算逻辑而是 watermark、迟到处理策略和读取路径上的时间差。我会在发布前加一道旁路对比把差异量化出来再决定是否可以放量。5.1 用旁路 diff 验证在线与离线一致性旁路对比的做法很直接新特征作业运行的同时把输出也写到一份独立 JSONL 文件中同时用同一份 YAML 配置触发的离线任务回算同一批窗口再按rider_id和window_start做对齐比较。下面是一个简化的比对脚本import json def load_features(path): result {} with open(path) as f: for line in f: r json.loads(line) result[(r[rider_id], r[win_start])] r[value] return result online load_features(online_features.jsonl) offline load_features(offline_features.jsonl) diff_cnt 0 for key, offline_value in offline.items(): online_value online.pop(key, None) if online_value is None: print(fmissing online: {key}) diff_cnt 1 continue if abs(online_value - offline_value) 0.01 * abs(offline_value) 1e-6: print(fdiff: {key} online{online_value:.3f} offline{offline_value:.3f}) diff_cnt 1 if online: print(fmissing offline: {list(online)[:5]}) print(fdiff count: {diff_cnt})这个脚本通过相对误差0.01加上绝对误差1e-6来容忍浮点计算的轻微差异。分母接近 0 的特征需要再过滤例如要求窗口内订单数至少为 5否则对稀疏样本做任何比较都没有意义。配送场景下我会按 10% 的骑手采样跑 30 分钟对比只要 diff 比例超过阈值就阻断发布。真正执行时不需要每天跑全量旁路对比只需要覆盖新特征首量阶段和容灾切换后观察 10 到 30 分钟即可。5.2 基于 Kafka 回放做特征变更后的冒烟验证离线对比通过后还差最后一步验证新特征在历史事件回放中是否稳定。通常会在一个独立集群上从 Kafka 某一位点开始重放事件并把回放结果落成文件随后再次执行上面的 diff。回放的任务参数可以这样启动flink run -d \ -D state.savepoints.dirhdfs:///flink/savepoint/feature_platform \ -c com.example.ReplayValidateJob \ feature-platform-1.0.jar \ --topic delivery_order_event \ --start-timestamp 1735689600000 \ --end-timestamp 1735789600000--start-timestamp和--end-timestamp使用毫秒级 epoch 时间回放范围建议覆盖一个完整高峰时段例如午间 10 点到 14 点。回放任务不写入 Redis只输出结果文件这样不会污染线上存储。回放过程中重点观察的是特征跳变次数同一个 rider 在相邻 1 分钟窗口内的特征值变化不应超过业务设定的合理幅度如果发生频繁跳变优先排查的是 watermark 延迟和迟到事件是否触发了过多重新计算。这套验证跑完后新特征才有资格进入线上灰度也才能在后续的模型复盘中做到在线离线同源可解释。本文还有配套的精品资源点击获取