ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

实时特征平台建设实践:从分钟级计算到稳定性保障

实时特征平台建设实践:从分钟级计算到稳定性保障 简介《美团配送实时特征平台建设实践》是一份针对实时计算与数据平台方向的技术分享PDF面向数据架构师、实时计算工程师及算法团队系统阐述配送场景下分钟级时效实时特征平台从0到1的建设路径解决烟囱式开发、重复建设与稳定性风险高的核心痛点。内容覆盖平台目标与整体架构、数据流处理、计算层设计、实时特征服务、稳定性建设及规模化演进重点介绍了订单-包裹-运单的数据建模、SQLUDF开发模式、拼图式合流解决乱序与端到端精确一次、基于内存计算的可扩展计算框架以及四层监控、三层降级、双机房容灾等稳定性实践。整份资源为1个PDF文件大小约56.23MB已有195人参与学习浏览。读者可从中借鉴美团在实时特征收敛、数据质量监控、查询性能优化和平台化架构升级方面的落地经验对规划自研实时特征平台或优化现有实时链路具有直接参考价值。1. 实时特征平台美团配送把特征从按天算压到分钟级的落地样本做实时特征的团队大多有过这种经历算法同学提需求时张口就是我要今天每小时的区域进单量一听业务没错但落到数据链路就成了黑匣子。数据从业务库同步到离线数仓要 T1Kafka 里的实时流又没人接最后算法只能硬着头皮用日级特征顶上去。这份《美团配送实时特征平台建设实践》讲的就是怎么把这层窗户纸捅破——用订单、包裹、运单三层建模刻画履约全程把特征生产收敛到一个分钟级时效、支撑调度/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 SQLUDF 开发模式与 DWD/DIM 分层从烟囱式到标准化美团把离线数仓的 SQLUDF 模式搬到了实时链路用 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:00SQL 里有两个参数值得关注。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[]数组代替ListFeature用原始类型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 画在一张表里数据一多系统性的资源瓶颈自然就浮现出来了。从那以后我每次接新的特征需求都会先问一句这个特征走哪条链路、兜底策略是什么、覆盖率指标挂在哪个面板上三个问题答清楚才动手开发希望帮到你。本文还有配套的精品资源点击获取
返回列表