简介:这是一份面向金融风控工程师、大数据开发及算法建模人员的智能风控在线特征系统实践分享。内容基于58同城2020年技术演讲,系统梳理了特征系统从离线到在线、从天级到秒级、从手动到自动的演进路径,并给出自然窗口、固定窗口、滑动窗口三类时间窗口及维度特征的基本概念,重点讲解Storm、Kafka Stream、Spark Streaming、Flink等实时计算框架的选型对比。针对滑动窗口计算、去重计算和字段提取三大难点,还介绍了延迟队列、顺序队列及自研TC框架的解决思路。整份资料共1个PDF文件,压缩包约1.69MB,图文结合地呈现背景、架构设计、特征生产、总结展望四大部分,适合希望系统理解智能风控特征架构的进阶读者。目前已有143人学习,可作为相关技术体系梳理与方案设计的参考。
1. 在线特征系统是什么:风控决策为什么绕不开它
看到“2-5+58”这个版本号,做过风控的朋友应该能猜到,这大概率是一份带着线上迭代编号的智能风控在线特征系统设计文档。版本号不重要,重要的是这套系统解决的根源问题:离线训练好的模型AUC有0.82,一上线日间决策就掉到0.76,排查半天发现不是模型问题,是线上喂给模型的特征和离线训练时对不上。在线特征系统,就是风控决策链路里“喂给规则和模型的实时特征生产工厂”,负责把埋点、订单、设备、行为等原始数据在毫秒级加工成规则和模型可用的特征。做反欺诈、贷前授信、交易拦截的策略和平台工程师,基本都躲不开这一环。这篇按“设计+实践”的顺序,把链路拆开讲清楚。
2. 在线特征系统的最小闭环:从埋点到特征读取怎么把链路串起来
2.1 选型理由:为什么风控在线特征不能照搬离线批处理
离线特征系统跑的是T+1的批处理任务:每天凌晨把前一天的数据从数据仓库捞出来,按实体ID聚合,算好特征写进特征表。模型训练时从这张表读特征,逻辑简单、口径统一。但线上实时决策不行,用户点一下借款按钮,风控系统必须在几百毫秒内返回能不能借、借多少。如果照搬离线逻辑,等Hive任务跑完再出特征,用户早走了。
所以在线特征系统第一个设计原则是:特征计算和特征读取分离。计算层用流式计算引擎做实时聚合,读取层用高速缓存存储已经算好的特征值,决策服务通过接口直接取数。两个层之间通过消息队列解耦,避免计算抖动直接打到决策链路上。
第二个选择是实时计算引擎的选型。我见过很多团队在Flink和Spark Streaming之间纠结。如果团队已经重度使用Spark,用Spark Streaming可以复用离线代码,但风控在线特征的典型场景是:高频窗口聚合、事件时间处理、迟到数据管理,这些正好是Flink的强项。我一般建议新项目直接上Flink SQL,理由很简单:特征逻辑多数是“近5分钟支付失败次数”“近1小时设备关联账户数”这类窗口聚合,Flink SQL写起来比DataStream API短一半,而且和离线Hive SQL语法接近,特征口径在两条链路上更容易保持一致。
第三是缓存选型。在线特征读取端几乎没有一个团队敢用MySQL扛高并发QPS,主流选择是Redis Cluster。特征是典型的read-heavy负载,写一次、读很多次,Redis的纯内存读性能足够。另一个容易被忽略的点:Redis key的TTL本身就是特征的生命周期管理工具,特征过期自动淘汰,比在应用层写定时清理任务省心得多。
2.2 实时计算层:用Flink SQL把窗口特征“算”出来
在线特征计算的核心是把“事件流”变成“特征值”。最常见的场景是:业务方上报行为事件,比如用户ID、设备ID、事件类型、金额、时间戳,系统需要实时统计每个用户近5分钟的支付失败次数、近1小时的设备关联账户数。用Flink SQL做这个事,代码比想象中短。
下面是一个从Kafka读取事件、按用户ID开1小时滚动窗口聚合支付失败次数的Flink SQL示例,这个模式能覆盖大部分风控窗口特征:
-- Flink SQL: 在线特征 - 支付失败次数窗口聚合 CREATE TABLE kafka_events ( user_id STRING, device_id STRING, event_type STRING, amount DECIMAL(10, 2), event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL '30' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'risk-events', 'properties.group.id' = 'feature-compute-group', 'scan.startup.mode' = 'latest-offset', 'format' = 'json' ); CREATE TABLE risk_features ( user_id STRING, feature_name STRING, feature_value DECIMAL(10, 4), window_end TIMESTAMP(3) ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://feature-meta:3306/risk_feature', 'table-name' = 'feature_snapshot' ); INSERT INTO risk_features SELECT user_id, 'pay_fail_cnt_1h' AS feature_name, COUNT(*) AS feature_value, HOP_END(event_time, INTERVAL '10' MINUTE, INTERVAL '1' HOUR) AS window_end FROM kafka_events WHERE event_type = 'PAY_FAIL' GROUP BY user_id, HOP(event_time, INTERVAL '10' MINUTE, INTERVAL '1' HOUR);这段SQL里有两个参数值得单独说明。WATERMARK FOR event_time定义事件时间的最大延迟容忍度为30秒,意思是超过当前事件时间30秒以上的迟到数据会被丢弃或视为过期。窗口函数用的是HOP(滑动窗口),窗口大小为1小时、滑动步长10分钟。这意味着同一时刻会同时计算6个窗口,每个事件进入6个窗口分别累加。滑动步长决定了特征更新的粒度:步长越短,特征越新鲜,但计算量和状态存储量线性上升。我们线上用的比较多的是10分钟步长,如果业务要求秒级更新,可以改成1分钟,但要做好状态膨胀的心理准备。
写入端我示例里用了JDBC连接MySQL,实际生产很少直接把特征写MySQL,更常见的做法是先写Kafka再异步刷到Redis。因为在线特征读取是高频路径,MySQL在几万QPS下会成为瓶颈,而Kafka天然具备削峰填谷的能力。这里保留JDBC只是为了说明“Flink SQL可以直接对接存储”,生产部署时建议换成Redis或HBase的connector。
2.3 特征读取层:给决策服务一个“拿得到、拿得快”的接口
计算层把特征算出来了,决策服务怎么拿?这是在线特征系统最容易拍脑袋的地方。我见过有的团队让决策服务直接查Flink的状态后端,或者直连HBase扫数据,结果就是特征接口RT(响应时间)从2毫秒抖到200毫秒,风控决策在一个大促瞬时流量下被打穿。
正确的做法是:特征读取层单独封装一个服务,对外只暴露一个批量查询接口,内部通过pipeline批量访问Redis。单个决策请求通常会携带用户ID、设备ID、订单ID等多个实体,一次决策需要几十个特征。如果循环调用Redis几十次,光网络开销就几十毫秒,必须用MGET或者pipeline一次取回。
// 特征读取服务: 批量取特征并处理缺失值 // keys 形如 "feat:user:1001:pay_fail_cnt_1h" func BatchGetFeature(ctx context.Context, keys []string) map[string]FeatureValue { // pipeline批量读取,避免N次RTT pipe := rdb.Pipeline() cmds := make([]*redis.StringCmd, len(keys)) for i, k := range keys { cmds[i] = pipe.Get(ctx, k) } _, _ = pipe.Exec(ctx) result := make(map[string]FeatureValue, len(keys)) for i, cmd := range cmds { val, err := cmd.Result() if err == redis.Nil { // 特征缺失: 走回源逻辑,而不是直接返回错误 result[keys[i]] = FeatureValue{Value: 0, Missing: true} continue } if err != nil { logger.Warnf("redis get failed, key=%s, err=%v", keys[i], err) continue } fv, _ := strconv.ParseFloat(val, 64) result[keys[i]] = FeatureValue{Value: fv, Missing: false} } return result }这段代码的要点有四个。第一,pipeline必须复用连接,不能每次请求新建连接,否则连接建立的开销比读特征还大。第二,特征缺失和特征值为0是两回事,缺失要打标记,让上层模型或规则决定是填默认值还是走拒绝策略,不能静默填0,不然“没数据”和“数据是0”在风控策略里含义完全不同。第三,Redis的Get命令返回的是字符串,在线特征服务里统一转成float64,但注意JSON序列化和反序列化的精度损失,这个坑后面单独讲。第四,这个接口必须是批量接口,单个查询接口再快也扛不住风控决策的并发量。
读取层的缓存策略也要设计。特征写入Redis之后,TTL设多久?太短,特征还没被读就过期;太长,用户行为变化了特征还是旧的。我们的经验是:按特征类型分开设置。强时效特征(如设备关联账户数)TTL设5分钟,弱时效特征(如用户注册天数)TTL设1小时,规则和模型都允许用“稍旧”的弱时效特征。避免一刀切,特征新鲜度对风控决策的影响比大多数人想象的大。
3. 离线在线一致性:怎么让同一份数据在两条链路上算出同一个数
3.1 为什么口径不一致是风控特征的“黑匣子”
在线特征系统上线后最诡异的一个现象是:模型在离线训练集上表现很好,上线后效果直线下滑,但你查不到任何代码bug。原因往往不在模型,而在特征口径。离线特征是用Hive在T+1批量算的,在线特征是用Flink实时算的,两条链路对“同一个特征”的定义有细微差别,比如时间窗口的边界、事件时间的取值、NULL的处理,累积起来就会让特征分布发生偏移。
举一个真实例子。离线特征“近7天支付失败次数”用的是自然日+自然小时的概念,即从当天0点开始算;而在线Flink任务用的是滚动窗口,从任务启动时刻开始算。若任务是上午10点启动的,则线上“近7天”实际是“7天前10:00到现在”,和离线“从昨天0点到现在”差了10个小时的数据切片。模型训练时见过的是“0点切”的分布,线上喂的是“10点切”的分布,特征值在边界处出现系统性偏差,AUC自然往下掉。
这就是在线特征系统里最典型的“黑匣子”:特征计算程序本身没写错,但时间语义错了。解决方向不是在程序里打补丁,而是把“特征口径”提升为一等公民来管理——每个特征除了代码实现,还必须有一个口径描述文档和执行验证脚本。Flink SQL和Hive SQL共用一套事件定义和窗口表达式是基本要求,落地时尽量用同一个SQL生成器或至少共享同一份特征口径配置文件。
3.2 校验方案:把离线结果和在线结果拉出来“对账”
口径不一致不是靠review代码能发现的,必须靠数据对账。我们的做法是:每天挑一个样本窗口,从离线特征表里取一批实体ID的特征值,再从在线特征系统里拉同一批实体同一特征的值,逐条对比,偏差超过阈值就报警。
对账不能只比均值,要看分布差异和逐条差异。均值相同不代表每条都对,有可能大值小值互相抵消。下面这个脚本是典型的对账工具,核心逻辑是拉两份结果,计算逐条绝对误差和误差分布。
# 离线在线特征一致性校验脚本 # 用法: python check_consistency.py --feature pay_fail_cnt_1h --date 2025-01-15 import argparse import pandas as pd import requests def load_offline(feature, date): # 从Hive/离线特征表读取, 这里省略SQL执行细节 # 返回 DataFrame: entity_id, feature_value return pd.read_csv(f"/data/offline_feature/{feature}/{date}.csv") def load_online(feature, entity_ids): # 调用在线特征读取接口, 批量取特征 url = "http://feature-server:8080/batch_get" resp = requests.post(url, json={ "feature": feature, "entity_ids": entity_ids[:5000] # 限制单次请求数量,防超时 }, timeout=5) return pd.DataFrame(resp.json()["data"]) def main(): parser = argparse.ArgumentParser() parser.add_argument("--feature", required=True) parser.add_argument("--date", required=True) args = parser.parse_args() offline = load_offline(args.feature, args.date) online = load_online(args.feature, offline["entity_id"].tolist()) merged = offline.merge(online, on="entity_id", suffixes=("_off", "_on")) merged["abs_err"] = (merged["value_off"] - merged["value_on"]).abs() merged["rel_err"] = merged["abs_err"] / (merged["value_off"].abs() + 1e-9) # 核心输出: 误差分布和毛刺记录 print(f"样本数: {len(merged)}") print(f"绝对误差均值: {merged['abs_err'].mean():.4f}") print(f"相对误差大于10%%的样本数: {(merged['rel_err'] > 0.1).sum()}") # 列出偏差最大的前20个实体, 用于人工排查 worst = merged.nlargest(20, "abs_err")[["entity_id", "value_off", "value_on", "abs_err"]] print(worst.to_markdown(index=False)) if __name__ == "__main__": main()这个脚本里有几个选型细节。第一,为什么不一次拉全部实体?在线特征接口一次查几千个key没问题,但查几十万个key会把Redis打满,对账任务要放在低峰期,且分批拉取。第二,误差判断用相对误差而不是绝对误差,因为不同特征的量纲差异很大,支付次数是几十的量级,金额是几千几万的量级,用绝对误差阈值没法统一。第三,对账频率不是越高越好。每天跑一次能发现大部分问题,但窗口边界错位这类问题只有跑批时刻能暴露,所以对账时间点要覆盖特征窗口的边界,比如0点、10点、下午4点。我见过不少团队只固定早上8点跑对账,结果在10点边界错位的特征永远发现不了。
3.3 一致性参数:窗口、迟到数据与精度怎么设
离线在线对账发现不一致之后,怎么修?大多数情况是调Flink SQL里的三个参数。
第一个是窗口类型。离线SQL通常用GROUP BY加固定时间条件,比如WHERE dt = '2025-01-15',这是典型的T+1全量窗口。在线要复现同一个口径,不能直接用滚动窗口,得用事件时间对齐。具体做法是把窗口终点对齐到自然日边界,即凌晨0点。Flink SQL里可以通过计算window_end = DATE_FORMAT(event_time, 'yyyy-MM-dd')来手动对齐,而不是依赖HOP或TUMBLE自动生成的窗口边界。
第二个是迟到数据。离线特征天然包含当天全部数据,不存在迟到问题;在线特征要等数据流实时到达,网络抖动或客户端重试都会造成数据延迟。如果WATERMARK设太短,比如10秒,很多正常延迟的数据被丢;设太长,比如10分钟,特征一直等数据,读取时拿到的可能是旧的聚合结果。我们线上的经验值是:业务方上报延迟的P99是30秒,那WATERMARK设60秒,allowedLateness设5分钟,超过允许范围再走旁路补偿。补偿逻辑一般是把迟到数据单独算一份增量特征,合并到主特征上,但这套做起来成本高,初期可以不做,只保证核心特征的迟到容错。
第三个是精度。离线Hive的DECIMAL默认精度和在线Flink的DECIMAL可能不一致,比如离线算出来的1.23456789,在线可能因为精度设置截断成1.23460000,看似差别微小,但风控模型对特征的微小抖动往往很敏感。我们遇到过浮点精度不一致导致策略命中率波动0.3个百分点的案例。解决方案是统一规范所有特征值的精度:离线Hive和在线Flink都用DECIMAL(20, 6),写入Redis时统一转成6位小数字符串,读取时解析成float64。注意float64本身也有精度上限,超过15位有效数字会出现“看着相等的两个数实际不相等”,特征值精度控制在6位小数是安全线。
4. 存储与容灾:Redis大Key、降级与回源怎么设计
4.1 特征缓存设计:key结构、TTL与序列化
在线特征数据写进Redis,第一个要设计的是key结构。我见过最省事的做法是直接feature_name:entity_id,比如pay_fail_cnt:1001。这方案在特征数少的时候没问题,但风控场景动辄几百个特征,一次决策要拉几十个key,MGET的key数量一多,网络耗时和Redis内存碎片都会上升。
我们的做法是分层设计key:场景前缀 + 实体类型 + 实体ID + 特征类别 + 特征名。比如feat:user:1001:pay:fail_cnt_1h,feat:device:abc123:relation:linked_account_cnt_1h。好处有两个:一是同一个实体的同类特征在key空间里是连续的,可以用scan按前缀批量拉取,减少MGET的key数量;二是TTL设置可以按前缀统一管理,比如feat:user:*下面的弱时效特征全都1小时过期,不用对每个key单独设置。
序列化方式也要选。常见选择是String存JSON字符串和Hash存字段。我们的经验是:能不用Hash就不用Hash。原因很直接,Hash在key数量大的时候会触发Redis内部编码转换,从ziplist变成hashtable,内存占用会明显上涨,而且Hash的单个字段TTL没法单独控制,不利于特征生命周期的精细管理。简单方案是每个特征一个String key,value就是特征值的字符串表示,Redis本身对短字符串的存储优化做得不错。当某个实体的特征数量特别多时,可以用pipeline批量写入,一次性把该实体所有特征写进Redis。
TTL策略这块有个反直觉的经验:不是所有特征都要设短TTL。频繁过期的特征会导致缓存命中率下降,决策服务不得不频繁回源,回源压力一大,在线特征服务就会从“快路径”退化成“慢路径”。我们的做法是把特征按实时性分成三档:实时特征TTL 5分钟,准实时特征TTL 30分钟,静态特征TTL 24小时。静态特征比如用户注册天数、历史逾期次数,这些值一天甚至一周不变,没必要让它和实时特征一样频繁过期。TTL的最终目标是让特征在“不失效”和“不过期”之间平衡,而不是一味追求短TTL。
4.2 降级与回源:本地缓存、Redis、HBase三级兜底
在线特征系统最怕的不是Redis慢,而是Redis挂。线上瞬时流量打进来,Redis主从切换或者网络抖动,特征接口RT从2毫秒涨到500毫秒,风控决策链路直接超时。如果特征服务没有降级方案,整个交易系统都会跟着抖。所以三级缓存架构是标配:本地缓存(Caffeine/GoCache)→ Redis → HBase回源。
本地缓存放最热门的特征,命中率通常能做到20%-30%。Redis是主存储,承载绝大多数读取。HBase是最后的兜底,存的是全量特征快照,从Flink实时写入一份,作为Redis不可用时的回源数据源。HBase的读延迟在几毫秒到几十毫秒,比Redis慢,但比全链路超时好得多。
// 特征服务三级降级逻辑: 本地缓存 -> Redis -> HBase public FeatureValue getFeature(String key) { // 第一级: 本地缓存, Caffeine, 过期时间5秒, 防止热点key打垮Redis FeatureValue v = localCache.getIfPresent(key); if (v != null) { return v; } // 第二级: Redis集群 try { String val = redisClient.get(key); if (val != null) { FeatureValue fv = parseFeature(val); localCache.put(key, fv); return fv; } // Redis有值但为空串, 可能是占位符, 直接返回缺失 if (val.length() == 0) { return FeatureValue.missing(); } } catch (Exception e) { // Redis超时或连接异常, 不抛异常, 走下一级 metrics.increment("redis_timeout_total"); } // 第三级: HBase回源 try { String val = hbaseClient.get(key); if (val != null) { FeatureValue fv = parseFeature(val); // 回源成功, 回填Redis, 带短暂TTL防止雪崩 redisClient.setex(key, 30, val); return fv; } } catch (Exception e) { metrics.increment("hbase_timeout_total"); } // 三级全部失败, 返回缺失值, 让上层策略决定 return FeatureValue.missing(); }这段降级代码有几个关键细节。第一,本地缓存过期时间必须远小于Redis的TTL,这里是5秒,作用只是削峰,不是为了兜底数据一致性。第二,Redis和HBase的异常不能抛出去,一旦抛出,决策服务那边会直接报错,导致本可以走降级的请求被拒。第三,HBase回源成功后要回填Redis,但TTL要设短,比如30秒,防止“回源风暴”把HBase打垮——如果Redis抖动恢复后大量key同时过期,所有请求都回源HBase,HBase也会被打爆。第四,指标埋点redis_timeout_total和hbase_timeout_total必须提前加,线上没有监控就讨论降级策略,等于在裸奔。
回源还有个容易漏的点:当Redis里key已过期删除,但HBase里也没有该特征时,要在HBase层返回明确的空标识,并在Redis写入一个占位符,避免每次请求都穿透到HBase。占位符的value可以约定为一个特殊标记,比如__MISSING__,读取端识别到这个标记就返回缺失值。这叫缓存穿透保护,不做的后果是HBase在特征缺失率高的场景下QPS暴涨,回源链路比Redis先挂。
4.3 容量规划与毛刺排查:Redis内存、连接数和慢命令
在线特征系统的容量规划,核心盯三个指标:Redis内存使用率、连接数、慢查询数。内存使用率达到70%以上就要警惕,Redis的内存淘汰策略默认是noeviction,内存满了直接返回OOM错误,特征写入失败导致缓存命中率下降,系统性能会断崖式下跌。建议内存水位超过60%就提前扩容或者清理冷特征。
连接数这块容易被忽略。特征服务如果是用连接池访问Redis,池大小要结合QPS和P99耗时来估算。一个常用的经验公式:连接池上限 ≈ QPS × P99延迟(秒)× 2。如果QPS是5000,P99延迟是3毫秒,那么连接池上限大约是5000×0.003×2=30,太高或太低都不合适。太高浪费资源,太低在高并发下会排队,RT飙升。
慢查询的排查方法很简单:开启Redis的SLOWLOG,设置阈值为5毫秒,定期拉取慢命令的key列表。如果发现某个key频繁出现在慢日志里,大概率是碰到了大Key问题。大Key指的是单个key的value过大,比如某个用户ID下积累了上万个特征字段的Hash,读这个Hash一次就要几十毫秒。解决办法是在写入端做拆分,把一个大Key拆成多个小Key,或者改用String存储并限制单特征value长度。大Key的排查和治理是Redis运维里最常踩的坑,后面那章会单独展开。
5. 生产环境踩坑排查:延迟、口径、抖动与容量
5.1 现象:线上特征比离线多算了一倍——时间窗口边界错位
线上某一个“近24小时支付失败次数”的特征,对账时发现比离线结果普遍高出近一倍。排查时先怀疑是Flink窗口开大了,但翻代码发现窗口大小和离线SQL完全一致。最后定位到问题是时间窗口的起点没有对齐自然日。离线SQL的“近24小时”是昨天0点到当前时刻,Flink滚动窗口是任务启动时刻起算的24小时滚动块,两边的数据切面差了整整一个相位。
原因是Flink的TUMBLE窗口是依据流数据的事件时间自动划分的,窗口边界是整数时间点(比如0点、1点),而不是“从任务部署开始算”。如果数据源事件时间本身被业务方填错,或者Kafka消息的事件时间戳用的是消息生产时间而不是业务发生时间,窗口计算就会基于错误的时间切分。
解决方法是把窗口对齐逻辑显式写出来。用OVER窗口或者手动计算window_start和window_end,确保窗口切分点与离线Hive的dt分区边界完全一致。修改后对账通过,线上特征恢复正常。这个案例给我们的教训是:在线特征和离线特征即使SQL语义相同,也要在窗口边界参数上做显式校验,不验证就上线等于埋雷。
5.2 现象:Redis大Key导致全链路RT毛刺
某天下午特征服务P99延迟从3毫秒突然涨到80毫秒,持续时间约十几秒,每10分钟出现一次,和特征写入任务的时间点高度吻合。排查Redis的慢日志,发现了几个单次读取超过40毫秒的HGETALL命令,对应的key是一个用户ID下挂了8000多个特征字段的Hash。
原因有两个。第一,某个渠道的爬虫或恶意注册行为,让一个设备ID关联了大量用户特征,特征拼接服务把该设备ID下所有关联用户ID的特征全写进了同一个Hash。第二,Hash在字段数超过512个时会从ziplist转成hashtable,读取时会阻塞Redis单线程。这两个因素叠加,形成了慢命令。
解决措施分两步。第一步是写入端拆分:把同一实体的特征按类别拆成多个key,限制单个key的字段数上限,Hash字段数超过100就自动拆成多个Hash。第二步是读取端改造:把对Hash的HGETALL改成按需取字段的HMGET,只取模型和策略真正会用到的字段。现实中很多团队不管需不需要,一股脑全取,才把简单读操作变成了慢操作。大Key治理没有捷径,唯一的办法就是写入端约束和读取端裁剪同时上。
5.3 现象:浮点精度不一致导致模型分数忽高忽低
风控模型同学反馈,同一批用户在同一天的模型评分,在线和离线差出5-8分,虽然规则没变、特征逻辑没变,但特征的输入值存在肉眼可见的微小差异。我们把线上特征值和离线特征值逐一打印出来,发现类似1234.560001和1234.559998的差异频繁出现。
原因是离线Hive用的是DOUBLE类型存储特征,在线Flink写入Redis时转成了DECIMAL(10,4),保留4位小数,读取时再解析成float64——精度被截断后再转回浮点,和原始的DOUBLE必然有微小误差。模型对每个特征的微小变化都要计算梯度,几个特征的误差叠加,评分波动就被放大了。
解决方式是把特征精度规范化做进整个链路:Flink SQL里计算特征时统一转成DECIMAL(20,6),Redis里存字符串时保留6位小数,读取解析时float64直接解析,不要先转float32再转float64。同时校验脚本里把相对误差阈值从“大于1%报警”调整成“大于0.1%报警”,因为精度问题的误差通常小于0.1%,但影响却可能不小。精度这个坑最讨厌的是它不报错,数据看着都对,就是结果差一点,属于典型的需要靠严格对账才能暴露的问题。
5.4 现象:回刷任务把在线特征缓存打爆
某次策略迭代需要把“近30天最大借款金额”这个特征从历史数据回刷,回刷任务直接把历史全量用户计算后的特征写进Redis,结果Redis内存使用率从40%暴涨到85%,触发大量key驱逐,在线特征命中率下降,全链路RT跟着上涨。
原因是我们把离线回刷任务的产出直接写进了在线Redis,没有区分“离线特征表”和“在线特征缓存”。回刷历史特征的正确姿势应该是:先把结果写入HBase全量特征表,再通过一个受控的预热任务,按线上实际访问的热度分批写Redis,而不是一把梭全量写。Redis缓存的价值在于保存热数据,全量写冷数据只会挤占热数据的空间。
解决措施是建立回刷任务分级机制:全量回刷只能写HBase,Redis只能由在线计算链路写入;如果确需预热Redis,必须分批执行,并控制写入速率,避免Redis的info memory水位超过50%才发警报。另外要在回刷任务里加一个最大写入条数限制,防止任务配置错误导致批量写爆。
6. 进阶:给在线特征系统上一道“后悔药”——特征回刷与版本回滚的实操
在线特征系统上线半年后,你会遇到一个离线特征没有的问题:特征口径要改。比如“近7天支付失败次数”原来只统计支付渠道的失败事件,现在要加上协议支付失败。直接改Flink任务,线上特征从改动时刻起就是新口径,但模型训练数据和历史特征还是旧口径,这会导致“新旧特征混用期”,模型分数分布突变。
给特征系统上“后悔药”的具体做法是把每个特征都带一个版本号。feat:user:1001:pay:fail_cnt_1h:v3,版本号从v1到v3。每次口径变更,新版本特征写入新key,旧版本key继续保留,决策服务通过配置中心动态切换读取哪个版本。这样做的成本是多一份Redis存储,但换来的是灰度切换和秒级回滚能力:如果v3有问题,把配置切回v2,线上特征立即恢复,不用等Flink任务重启。
回刷的思路也依赖版本号。改口径后,先用离线任务重算历史一段时间的特征,写进v3的HBase快照,再分批预热到Redis。预热顺序按线上特征访问热度排,先刷高频实体,再刷中频实体。同时从切换配置那一刻起启动对账任务,持续1周观察v3和v2的特征分布差异,一旦发现偏差超过预期,就在配置中心回滚到v2。
我自己的习惯是,每年做一次在线特征系统的“退役演练”:挑一个低频特征,停写、停读、观察一周,确认没有告警、没有模型效果波动,再清理代码和存储。这个习惯已经帮我提前暴露过两次问题,一次是特征还在被旧的策略规则引用但代码里已找不到出处,另一次是Redis里积累了上百万的孤儿key。特征治理不是临上线前做一次,而是周期性动作。
每一版特征从设计到部署,都要在这种细节里反复掐一遍。希望这篇文章能让你少踩几个我已经趟过的坑,后面做在线特征系统时,心里更有底。
本文还有配套的精品资源,点击获取