简介:这是一份面向金融科技、大数据风控领域工程师与架构师的技术分享PDF,内容来自58同城资深数据开发工程师的实践演讲,系统梳理智能风控在线特征系统的设计思路与落地路径。资料从2017年网络黑产背景切入,讲解特征系统与规则、模型策略的协作关系,并围绕时间窗口特征、维度特征、滑动窗口计算等核心难点展开,对比Storm、Kafka Stream、Spark Streaming与自研TC框架的优劣,给出延迟队列、顺序队列的具体解法。包体为单个PDF文件,大小1.69MB,便于按章节精读与检索。文中还展示了数据中心、计算中心的流批一体架构及数据字典设计,适合希望提升实时特征计算能力、优化风控模型效率的读者参考。已有143人学习下载,内容兼具背景认知与工程实践细节,对理解在线特征生产链路和实时计算选型有直接帮助。
1. 智能风控在线特征系统:为什么特征上线比模型上线更让团队睡不着
风控系统的效果翻车,十次里有七次不是模型不行,而是特征没送到位。我见过一个信贷团队,模型离线 AUC 提升了 0.03,上线之后逾期率反向涨了五个点,最后定位到根因:线上实时算出的“近 7 天申请次数”和离线训练表里的口径差了 12 小时的数据。这份《2-5+58智能风控在线特征系统设计与实践》的内部代码号,就是用来解决这类问题的:把离线批量算好的特征,变成在线实时可查、口径一致、延迟可控的服务。本文按该系统最常见的落地路径展开,覆盖架构选型、特征在线化改造、发布验证和踩坑,适合风控算法、后台开发和数据平台三类读者照着搭一套最小可用版本。
2. 在线特征系统的架构选型:三级计算与存储分层的取舍
2.1 特征从离线到在线要跨的三道坎
风控特征有一个天然矛盾:离线训练时你用尽全量历史数据,特征越全越好;线上推断时你只有几十毫秒,不可能现场扫一遍 Hive 表。在线特征系统要做的就是把“能离线算的先算好,不能离线算的想办法实时算”,并且保证两边口径一致。
第一道坎是延迟。离线特征 T+1 生成完全没问题,线上风控要求特征读取加计算整体在 50 到 100 毫秒内返回。第二道坎是口径。离线用 Hive SQL 算出的 COUNT、SUM、AVG,到了实时计算要用 Flink 或 Spark Streaming 重写一遍,稍不注意窗口边界、时间字段单位、空值处理就会产生偏差,模型拿到的特征和训练时对不上。第三道坎是数据量。把所有特征一股脑塞进 Redis 不现实,全量用户特征可能上亿 key,成本高且命中率低,必须按特征的使用频率和更新频率分级存储。
2.2 离线/近线/在线三级计算的分工
我在实际项目里习惯把特征按更新时效分成三级,对应三条计算链路。离线级(T+1)处理变化慢的强变量,比如用户注册时长、历史逾期次数、设备首次出现日期、多头借贷历史。这类特征用 Hive 或 Spark 批处理,每天凌晨算好写入特征表,再同步到在线存储。近线级(分钟到小时级)处理中频变量,比如近 1 小时申请次数、IP 段聚集度、同设备关联账号数,用 Flink 或 Spark Streaming 做微批,延迟目标设在 1 到 5 分钟。在线级(毫秒级)处理只跟本次请求相关的变量,比如当前请求 IP 是否命中黑名单、设备指纹是否异常、本次输入的手机号段、实时埋点行为,这类必须请求进来现算。
三个级别缺一不可:全部放线上一是算不过来,二是很多特征(比如历史逾期次数)在线根本拿不到全量数据;全部放离线又满足不了实时拦截的需求。我的经验是一个特征从“离线能算”升级为“在线能算”前,先看它的变化频率和风控价值,只有两个都达标才值得迁移。
| 计算级别 | 更新时效 | 典型特征 | 存储介质 | 失败影响 |
|---|---|---|---|---|
| 离线 | T+1 | 历史逾期次数、注册时长、多头借贷 | Hive -> Redis | 特征陈旧,可用兜底值 |
| 近线 | 分钟~小时 | 近1小时申请次数、同设备关联数 | Flink -> Redis | 短时不准,影响评分 |
| 在线 | 毫秒 | 当前IP黑名单、设备指纹、实时行为 | 本地缓存 + Redis | 直接拦截失败,需降级 |
2.3 在线特征服务的存储选型与 Key 设计
在线特征服务的存储层,我见过三套方案:纯 Redis、Redis 加本地缓存、HBase 加 Redis。风控特征大部分是“一个实体 ID 对应一组特征值”,读多写少、单个 value 不大,Redis 的 Hash 结构最合适。HBase 适合特征维度极多、需要按列族批量扫描的场景,但运维成本高,在线延迟也不如 Redis 稳定。本地缓存适合高频热点特征,比如黑名单 IP、设备风险等级,命中率极高且允许分钟级更新延迟。
一个常被忽略的细节是 Key 设计。风控特征最少需要三个维度:特征名、实体 ID、特征版本。我常用的格式是featureName#entityId#version,比如apply_cnt_7d#U1000234#v20240501。版本号必须进 Key,否则特征口径调整后,旧缓存数据未过期,线上会读到新旧混杂的脏数据。TTL 根据特征更新频率定:T+1 特征设 24 到 48 小时,近线特征设 2 到 4 小时,在线特征通常不设 TTL,靠本地缓存淘汰。
public class FeatureFetcher { private LoadingCache<String, FeatureValue> localCache; // L1 本地缓存 private RedisTemplate<String, String> redis; // L2 Redis public FeatureValue getFeature(String featureName, String entityId, String version) { String key = featureName + "#" + entityId + "#" + version; FeatureValue v = localCache.getIfPresent(key); if (v != null) { return v; } String raw = redis.opsForValue().get(key); if (raw == null) { return null; // 上层走特征缺失兜底逻辑 } FeatureValue parsed = parse(raw); localCache.put(key, parsed); // 回填本地缓存 return parsed; } }这段代码没什么高深之处,但三个参数必须说明:version不传默认值的后果我之前已经说了;localCache我建议配置 maximumSize 在 10 万到 50 万之间,过大导致 GC 压力,过小命中率上不去;raw为空时不要当场计算,而是返回 null 让上层走兜底,否则每个请求都在线算一遍,等于把保护层拆了。另外,Redis 操作务必设置超时,我一般设 10 到 20 毫秒,超时直接返回 null,不能因为特征系统拖垮主流程。
3. 核心特征在线化改造:时间窗、序列特征与外部数据的实时计算
3.1 把 Hive SQL 的时间窗聚合改写成 Flink 实时任务
时间窗聚合是风控特征里最常用也最容易出问题的一类。离线口径是“近 7 天申请次数”,Hive 里一句count(*) group by uid就完事;到线上实时计算,你要决定两件事:窗口是滚动还是滑动,事件时间还是处理时间。
我先给一个离线版 SQL 做参照:
-- 离线特征表:每天跑一次,dt 是分区字段 SELECT uid, COUNT(*) AS apply_cnt_7d, AVG(amount) AS apply_avg_amount_7d FROM apply_record WHERE dt >= date_sub('${today}', 7) AND dt <= '${today}' GROUP BY uid;这段 SQL 看着简单,但有一个隐藏坑:dt是申请日期,不是业务发生时间。如果一笔申请跨天落库,dt按入库日期算,离线口径会对不上线上实时的事件时间。这也是为什么很多人离线特征和在线特征天然不一致的原因之一。
改写为 Flink SQL 近线任务,常见做法是开一个 7 天滚动窗口,按事件时间聚合:
INSERT INTO online_feature_sink SELECT uid, COUNT(*) AS apply_cnt_7d, AVG(amount) AS apply_avg_amount_7d FROM apply_record_stream GROUP BY TUMBLE(event_time, INTERVAL '7' DAY), uid;这里最关键的参数是 Flink 的状态过期时间。滚动窗口本身会在窗口结束时清理状态,但如果你用的是 OVER 窗口或者需要保留中间结果,状态会无限增长。我一般会给算子设置state.ttl,例如table.exec.state.ttl配置为 2 天,防止 7 天窗口的中间状态把堆内存撑爆。另一个参数是水位线:事件时间模式下,乱序数据会导致窗口提前关闭,建议允许 1 分钟内迟到数据,再晚的直接丢弃并记录指标,用于监控数据质量。
3.2 IP 与设备风险特征的在线打标
IP 黑名单、设备风险等级、手机号段风险这些特征,不适合批量算完再存,因为请求进来之前你根本不知道下一个 IP 是谁。常见做法是分两层:第一层用布隆过滤器做粗筛,判断“这个 IP 是否在历史风险集合里”;如果命中,再走一次精确查询拿到详细标签。
布隆过滤器有一个误判率参数要设定:fpp(false positive probability)设 0.01 时容量 1000 万的集合只需要约 1200 MB 内存,设 0.001 时约 1700 MB。风控场景误判会导致正常用户被标记为风险,所以我建议 fpp 设 0.001 而不是 0.01,代价是多几百 MB 内存,换取更低的误伤。同时布隆过滤器不支持删除,黑名单 IP 到期移除时,要重建过滤器,否则过期风险 IP 永远占着内存。我之前踩过这个坑,后来改成每天凌晨从离线风险库重建一次布隆过滤器,运行期只添加不删除。
3.3 特征配置化:把特征上线从发版变成改配置
在线特征系统如果每个特征都要改代码发版,运营成本会随特征数量线性增长。更常见的做法是做一套特征配置中心,把特征的定义、计算逻辑、存储位置、TTL、兜底值都放在配置里,由配置中心下发到特征服务动态加载。
一个最小配置结构长这样:
{ "featureName": "apply_cnt_7d", "version": "v20240501", "level": "neartime", "computeType": "flink_sql", "computeLogic": "SELECT uid, COUNT(*) FROM apply_record_stream GROUP BY TUMBLE(event_time, INTERVAL '7' DAY), uid", "storeKeyTemplate": "apply_cnt_7d#{uid}#{version}", "ttlSeconds": 14400, "fallbackValue": 0 }这里的fallbackValue最容易被人忽视。风控特征查询失败时不能直接让请求报错,必须给一个默认值。默认值怎么定是门学问:apply_cnt_7d这种“次数越高风险越高”的特征,兜底设 0 意味着“更信任这个用户”,如果缓存故障发生在高风险时段,反而会放行;更稳妥的做法是设成历史分布的 P50 或 P90 值,宁可误杀不放过。我把这个思想叫“特征兜底偏向风控”,不是所有特征都适合给 0 或 null。
4. 特征发布与一致性校验:从离线表到在线缓存怎么不出错
4.1 离线特征表同步到 Redis 的链路
每天凌晨批处理产出特征快照后,需要同步到在线存储。这里推荐用分片同步而不是单线程全量灌数据。假设离线特征表有 1 亿行,按 uid 哈希分成 64 个分片,每个分片一个同步任务,整体同步时间能从 1 小时压到 15 分钟以内。
一个典型的数据同步任务骨架:
#!/bin/bash # 从 Hive 导出当日特征快照,按 uid 哈希分片 hive -e " SET hive.exec.dynamic.partition=true; INSERT OVERWRITE TABLE feature_snapshot_daily PARTITION(dt='${TODAY}') SELECT uid, feature_name, feature_value FROM feature_offline WHERE dt='${TODAY}' DISTRIBUTE BY HASH(uid, 64); " # 逐分片同步到 Redis for shard in $(seq 0 63); do hadoop fs -cat /warehouse/feature_snapshot_daily/dt=${TODAY}/shard=${shard}/* | redis-cli --pipe -h ${REDIS_HOST} -p ${REDIS_PORT} & done wait这个脚本有两个参数值得说明:DISTRIBUTE BY HASH(uid, 64)是保证同一个 uid 的所有特征落在同一个分片,这样 Redis 写入时天然避免跨分片事务问题;redis-cli --pipe是批量导入模式,比逐条SET快一个数量级,但要注意一次管道的数据量,建议每 10 万条一批,避免 Redis 响应缓冲区积压导致连接超时。
4.2 影子验证:离线回放与在线比对
特征同步上线最怕的不是慢,而是“悄悄不一致”。我见过一个特征在离线表里用了unix_timestamp存时间,在线实时计算里用了字符串时间,两边没 diff 出来,模型上线后逾期率涨了三个点才定位到。所以特征发布前必须做影子验证:用真实历史流量回放,比对离线值和在线值。
回放比对的逻辑很简单:取最近 7 天的线上请求日志,逐条根据请求里的 uid、ip、device_id 去查在线特征服务,再把结果和当天离线特征表里对应记录做差。一致率低于 99.5% 就禁止该特征放量。注意 99.5% 是整体容忍度,单条特征的差异要按错误类型分类:小数值浮点差、窗口边界差、空值差,分别设独立的容忍度。
# 离线值 vs 在线值比对脚本 import json confident_count = 0 total_count = 0 diff_examples = [] with open("replay_result.jsonl") as f: for line in f: record = json.loads(line) total_count += 1 offline_val = record["offline_value"] online_val = record["online_value"] if abs(offline_val - online_val) < 0.001: confident_count += 1 else: diff_examples.append(record) if len(diff_examples) >= 10: break print(f"一致率: {confident_count / total_count:.4f}") print(f"差异样例: {json.dumps(diff_examples, indent=2)}")这段脚本里我刻意把0.001作为浮点容忍度,这是血泪教训:特征值经常是概率、比率,离线用 Decimal 在线用 Double,正常计算产生的浮点误差在 1e-6 量级;如果差到 0.001 以上,基本可以断定是口径问题而不是精度问题。跑完脚本后还要人工看差异样例,光看一致率数字不够,我遇到过一致率 99.9% 的情况下,差异的那 0.1% 恰恰是风险最高的头部用户。
4.3 灰度发布与特征开关
特征系统改动不能一把梭全量切换,必须灰度。常见的做法是在特征服务里嵌入“特征开关”逻辑:请求进来先读配置中心,判断当前请求的 uid 是否命中灰度分桶。分桶规则用hash(uid) % 100,灰度比例从 1% 开始,逐步调到 5%、 20%、 50%、 100%。灰度期间实时对比灰度和非灰度的特征命中率、平均耗时、兜底触发率,任何一项异常立即把开关调回 0%。
特征开关的配置项要包含三个字段:特征名、生效比例、生效版本。比如apply_cnt_7d当前线上是v20240501,要发布v20240601,配置中心先下发{"featureName": "apply_cnt_7d", "version": "v20240601", "ratio": 5},特征服务收到后只对 5% 的请求用新版本,其余 95% 仍走旧版本。这个机制的好处是回滚不需要改代码,只需要把 ratio 调成 0,天然给团队留了后悔药。我在实际项目中,特征灰度发布通常比模型灰度多观察一个完整业务周期(比如 7 天),因为很多欺诈特征在日维度上才有区分度。
5. 在线特征系统常见问题与排查:五个线上坑的血泪经验
5.1 缓存穿透:风控大促瞬间打到下游数据库
现象:大促期间风控 QPS 冲到平时的 20 倍,特征服务 Redis 命中率从 90% 掉到 40%,数据库连接数被打满,特征查询平均延迟从 20ms 飙到 2 秒。
原因:活动期间大量新用户涌入,这些用户没有历史特征,离线表里自然查不到,请求全部穿透 Redis 打到数据库。数据库没有这么高的读能力,连接池耗尽后所有查询排队。
解决:为每个特征加一个“空值占位符”。查询 Redis 时如果 key 不存在,写一个 TTL 为 30 到 60 秒的空值占位,后续相同请求直接命中占位符,不再穿透。同时限制单 IP 每秒最大特征查询数,超限直接拒绝并返回兜底值。这个坑的特征在于:你很难从监控里发现“穿透”本身,只能看到延迟突增,排查时要先看 Redis 命中率曲线。
5.2 特征值延迟到达:上游数据没到,下游算了半截
现象:某个近线特征在每天凌晨 2 点会有 10 分钟的值断层,期间返回的全是兜底值;白天偶尔也会出现某分钟窗口特征值为 0。
原因:上游日志数据因为 Kafka 消费 lag 或 HDFS 小文件问题延迟入库,Flink 窗口已经触发计算,但数据还没到达,导致窗口聚合结果偏小甚至为 0。离线表第二天修正后,在线特征已经用错误值参与了一天的风控决策。
解决:给 Flink 任务设置“允许迟到数据”的时长,我常用的是窗口结束后再等 1 分钟,迟到数据触发一次增量更新。同时特征服务端增加“特征新鲜度”监控:对于近线特征,如果当前时间减去特征最后更新时间超过阈值,直接标记该特征不可用,返回兜底值而不是返回旧的正常数值。
5.3 数据倾斜与热点 Key:某个头部用户的特征被反复查询
现象:Redis 某个节点的 CPU 使用率比其他节点高 5 倍,特征服务整体延迟正常,但错误率上升。
原因:头部欺诈分子可能同时发起大量申请,同一个 uid 的特征被高频查询,Redis 单分片热 key 出现。更隐晦的是特征表里某个 IP 段覆盖了大批量用户,导致按 IP 聚类的特征在同一时间段内被密集访问。
解决:热 key 加本地缓存,我前面说的 L1 本地缓存对这种情况特别有效,命中一次就在本地缓存中保留,后续请求不再打 Redis。热 key 本身需要在写入时做散列拆分,比如把ip_risk#1.2.3.0/24拆成ip_risk#1.2.3#shard0到ip_risk#1.2.3#shard15共 16 个分片,查询时随机选一个分片读取。拆分数量不宜过多,否则维护成本大于收益。
5.4 模型用的特征和线上算的特征不是同一套
现象:新模型上线后效果低于离线评估,排查了半天发现线上特征服务里apply_amount_avg_7d还在用旧口径(只算成功申请),而新模型训练时用的新口径(包含拒绝申请)。
原因:特征更新了,但模型版本和特征版本没有绑定。特征服务的配置信息里没有记录“这个特征属于哪个模型版本”,导致新模型上线时用了旧特征。
解决:在做特征服务时就把“模型-特征版本映射表”建好,每次模型上线前自动校验:模型依赖的特征版本是否全部在线可用。校验不通过直接拒绝发布。这个映射表可以是一个简单的配置文件,发布系统读取后自动比对,不允许人工跳过。
5.5 监控盲区:只看接口成功率,不看特征覆盖率
现象:特征服务接口成功率 99.99%,但风控模型的效果在缓慢下降,每周逾期率都在涨,但没人察觉。
原因:特征覆盖率(即请求能拿到非兜底特征值的比例)从 95% 慢慢掉到了 85%,兜底值使用比例升高,模型实际输入和质量变差。接口成功率只反映“请求有没有返回”,不反映“返回的值可不可信”。
解决:必须为每个特征单独埋“值来源”指标,区分三类:命中缓存、空值占位、兜底值。按特征、按接口维度统计兜底比例,超过阈值就告警。我在这个坑上吃过亏,后来把“兜底率”列为线上 dashboard 的前三个指标之一,和成功率并列。“接口正常但效果变差”这类慢故障,大多数情况下是特征覆盖率和新鲜度出了问题。
6. 最后一道关:用一份检查清单验收你的在线特征系统
在线特征系统跑起来只是开始,验收才是决定团队能不能睡得着觉的环节。我每接手一个风控特征系统,都会按固定清单过一遍,不发版也要每月自查。这份清单的核心是五个数字:特征覆盖率、特征新鲜度、查询成功率、P99 延迟、离线在线一致率,它们分别对应“特征有没有、数据新不新、服务稳不稳、延迟够不够、口径对不对”。一个正常的系统至少满足:覆盖率大于 99%,新鲜度延迟低于 5 分钟,成功率大于 99.99%,P99 延迟小于 100ms,一致率大于 99.5%。如果某一项不达标,对应的排查方向分别是:数据同步链路、Flink 窗口延迟、缓存穿透与热 key、特征序列化和网络开销、离线在线口径差异。
这里给你一个可以直接用的验收表格模板,我每次系统改造完就照着逐项打勾,曾在两周内发现 6 个潜在隐患:
| 验收项 | 通过标准 | 失败排查方向 |
|---|---|---|
| 特征覆盖率 | > 99% | 离线特征表是否有大量 key 未同步 |
| 特征新鲜度 | 近线 < 5min,离线 < 24h | Flink 消费 lag,同步任务失败告警 |
| 查询成功率 | > 99.99% | Redis 连接池、本地缓存淘汰策略 |
| P99 延迟 | < 100ms | 特征序列化格式、Redis 慢查询 |
| 离线在线一致率 | > 99.5% | 时间口径、空值处理、浮点容差 |
最后一个技巧和序列化有关:在线特征值我建议统一使用二进制格式(如 protobuf 或 MessagePack),不要直接存 JSON 字符串。JSON 的好处是调试方便,但序列化和反序列化开销在 P99 延迟上至少多消耗 5 到 10ms。批量特征读取时,用 pipeline 或 mget 而不是循环单条查询,Redis 往返次数从 N 次降到 1 次,P99 延迟往往能直接砍掉 30%。
我自己的习惯是每次台风控特征上线前,把这份清单贴在发布群里,没人敢免检。干这行越久越相信一件事:模型和算法只决定上限,特征系统才是底限;宁可模型笨一点,也不能让特征停下来。希望帮到你。
本文还有配套的精品资源,点击获取