ARTICLE DETAIL

资讯详情

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

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

实时特征平台架构:美团配送的分钟级统一与Flink动态计算实践 简介一份美团配送实时特征平台建设实践的技术分享PDF面向大数据实时计算开发者、平台架构师与算法工程同学。内容紧扣配送业务分钟级实时特征需求系统梳理从平台目标、整体架构到稳定性建设、规模化的完整演进路径。包内共1个PDF文件压缩包约56.23MB已有195人学习下载。资料重点讲解数据输入、加工、计算、输出四层架构通过SQLUDF模式提升效率以拼图式填充和上游合流应对数据乱序、实现端到端不丢不重计算层采用基于内存计算、无状态可扩展的升级框架支撑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 与 SQLUDF 开发模式3.1 为什么没有直接用 Storm 或 Flink 做特征计算2017 年时美团内部实时计算的主流选择是 Storm 和 Flink/Spark Streaming。但配送特征计算场景有几个特殊性特征计算逻辑密集且迭代极快算法同学经常要改一个 UDF 就发版验证特征量从几个涨到几十个每个特征都要消费宽表里若干列同时业务对计算耗时要求苛刻——每 10 分钟要处理的运单数据是千万级的。当时行业方案的短板恰好都踩在这几个点上方案优势在配送特征场景的短板Storm低延迟、原语丰富开发运维成本高SQL 化难度大状态管理能力弱Flink/Spark Streaming监控运维成熟、checkpoint 稳定基于关系数据库的计算模型扩展能力偏弱并发加不上去基于 RPC 的自研计算内存计算、无状态、可水平扩展需要自建调度和容错初期成本高最终选择了第三条路自研 FCSFeature Compute Service计算框架。核心设计是基于内存计算、计算无状态、可扩展。宽表数据加载到每个 Worker 节点的内存中计算任务通过 MQ 分发Worker 不保存跨任务的共享状态扩容就是加机器。这个思路的代价是放弃了 Flink 那样的统一状态管理换来的是计算耗时波动极小的确定性。3.2 SQLUDF借鉴离线数仓的开发模式实时特征平台真正提升研发效率的关键是把离线数仓的 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 层SQLUDF 模式不是万能的。对于需要迭代式计算的复杂特征比如基于骑手轨迹聚类出的驻留点识别仍然用 Java 原生算子写成独立计算任务通过 MQ 写回特征存储。平台统一提供 UDF 注册、版本管理、发布审批但引擎执行层对两种模式一视同仁。3.3 数据倾斜治理提前分片与能者多劳实时特征计算的数据倾斜比离线更头疼因为不能在发现倾斜后再做二次 MapReduce。FCS 的思路是在任务调度阶段就按区域分片。配送业务天然有区域属性一个商圈的运单量远大于另一个商圈。提前把宽表数据按 area_id 分片每个分片对应一组 FCS Worker分区内部再用 MQ 任务队列做能者多劳——哪个 Worker 处理完当前批次就从队列里取下一个任务避免固定分片导致的忙闲不均。定时任务 - FCS Worker(area1) - H2 索引 - MQ 定时任务 - FCS Worker(area2) - H2 索引 - MQ 定时任务 - FCS Worker(area3) - 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先查热路径上有没有不必要的对象分配再考虑调整堆大小。本文还有配套的精品资源点击获取
返回列表