简介:《智能风控在线特征系统设计与实践》是一份聚焦金融风控实时特征生产的专业分享资料,源自同城大数据应用实践,适合从事风控算法、数据开发与实时计算方向的工程师阅读。内容从2017年网络黑产规模切入,解释智能风控为何需要特征系统;随后系统梳理自然窗口、固定窗口、滑动窗口及维度特征等概念,并展示架构从离线到在线、从天级到秒级的演进。围绕实时计算的核心难点,材料重点讲解滑动窗口的延迟队列与顺序队列解法、去重计算的明细存储策略、字段提取与数据字典设计,同时对比Storm、Kafka Stream、Spark Streaming、Flink和自研TC框架的选型优劣。资源为单个PDF文件,共1.69MB,便于在电脑或移动端阅读。目前已有143人学习,适合需要深入理解风控特征生产链路、实时框架选型与工程落地细节的中高级技术从业者。
1. 智能风控在线特征系统设计与实践:一份来自 58 同城一线数据工程师的实战拆解
做风控的人都知道,特征就是命根子。2017 年国内黑产从业人员就超过 150 万,互联网上有将近 40% 的信息是虚假信息,年产值到了千亿规模。这意味着什么?意味着平台上的每一次内容发布、每一次交易行为,背后都可能有黑产在批量操作。规则策略可以通过用户行为特征定义阈值来命中,模型策略依赖用户行为、表征做综合判断,而所有这些策略的输入,都来自一个能在 50ms 内完成计算的特征系统。李文学在 2020 年分享的这份《智能风控在线特征系统设计与实践》,讲的就是 58 同城如何把离线特征计算搬到在线,实现天级到秒级的演进。这份 PDF 适合正在做风控特征平台、实时数仓,或准备自研实时计算框架的工程师,里面关于滑动窗口、去重计算、自研 TC 框架的取舍逻辑,比看十篇框架对比文章都来得实在。
2. 从离线到在线,特征系统的架构演进与技术选型逻辑
2.1 四个演进阶段:离线特征线上应用、实时流引入、自动化、全面支持算法
文本里的演进路径写得非常清楚,四个阶段分别对应不同的业务痛点和工程目标。第一阶段是离线特征线上应用,把 Hive 里算好的特征表同步到在线存储,供模型和规则查询。这个阶段的问题很明显,特征天级更新,黑产早就干完一波活了特征还没刷新。第二阶段从天级到秒级,引入实时流数据,Kafka 消息进来后实时计算窗口特征。第三阶段从手动到自动,把人工配置特征、人工上线特征的过程抽象成元数据驱动,节省数据开发效率。第四阶段从单一到全面,特征系统不只要支撑规则,还要支撑模型训练和在线推理,特征覆盖用户、设备、IP、银行卡等多个维度。
这里有一个容易被忽略的细节:离线到在线的切换不是简单的把计算挪个地方,而是数据模型、存储选型、计算语义全部要变。离线特征可以全量扫描、可以回溯、可以一天跑一次,在线特征必须 50ms 返回,窗口计算必须在内存中完成,还不能丢数据、不能重复计算。所以架构演进的核心不是把离线任务改成实时任务,而是重新设计一套面向在线场景的特征生产链路。
2.2 主流实时计算框架对比:为什么 Spark Streaming 会翻车
PDF 里给了一张 Storm、Kafka Stream、Spark Streaming、Flink 的对比表。Storm 的问题是延迟能到毫秒级但 Exactly-once 不支持、状态管理不完善;Kafka Stream 胜在轻量,但复杂窗口计算和批流融合能力弱;Spark Streaming 用微批模型,延迟在秒级,状态管理靠外部的有状态算子,复杂事件时间窗口支持有限;Flink 当时在 58 同城还处于引入阶段,一些内部组件和运维体系没有完全跑通。
关键结论是:如果直接用成熟框架硬扛特征计算的滑动窗口,会出问题。以 Spark Streaming 为例,它本质是微批处理,窗口边界按 Batch 对齐。当你需要 10:00 到 11:00 的滑动窗口特征,但窗口边界和数据事件时间的真实边界对不齐时,Spark Streaming 给出的结果会有系统性偏差。更大的问题是,长窗口(比如 24 小时甚至 7 天)在 Spark Streaming 里需要维护大量中间状态,一旦 executor 宕机,状态恢复就要从 checkpoint 重放,恢复时间不可控。这就是为什么 58 同城最终选择自研 TC(Time Calculator)框架——用延迟队列解决窗口过期,用顺序队列维护窗口内明细,用累加器、对比器、集合这些基础数据结构组合出精确的滑动窗口计算能力。
提示:选型时不要只看框架的官方文档写了什么,要拿自己的业务场景去压测。特征计算对精确性和延迟的要求,和普通实时 ETL 完全不是一个量级。
2.3 特征系统整体架构:数据中心、计算中心、统一服务三层
特征系统整体架构分为数据中心和计算中心,外加统一服务层。数据中心作为数据的统一出入口,采用流批一体结构,底下是实时数据仓库和离线数据仓库。实时数仓处理 Kafka 流式数据,离线数仓处理 Hive 表数据,两套数据通过统一的数据字典服务做元数据对齐。PDF 里专门问了一个问题:如果 Hive 没有 MetaStore 会怎样?答案是消息队列没有元数据管理,流批一体就无从谈起。数据字典的核心是让每条消息、每张表都有统一的 schema 登记,消费端和生产端按同一套元数据协议解析数据。
计算中心承载的是特征计算的执行引擎,包括离线计算引擎(Spark/Hive 批任务)和在线计算引擎(TC 框架处理实时流)。统一服务层对外提供特征查询 API,风控策略引擎在收到请求后,在 50ms 内并发获取多个特征值并合并结果。这里要重点理解:特征系统不是简单地把计算做成实时,而是把特征的计算结果实时化、服务化,让上游策略引擎可以像查数据库一样快速拿特征。
3. 特征生产的核心挑战:滑动窗口、去重计算与字段提取
3.1 滑动窗口的精确计算:延迟队列与顺序队列的设计
特征计算里最麻烦的是时间窗口特征。PDF 里把时间窗口分成自然窗口(比如 0:00 到 11:00,固定边界)、固定窗口(比如 9:30 到 10:30,固定时长)、滑动窗口(比如最近 1 小时,随当前时间持续滑动)和 Session 窗口(比如 30 分钟无操作断开)。自然窗口和固定窗口的边界是确定的,计算相对简单;滑动窗口的边界是持续变化的,每来一条新数据就要计算一次当前窗口内的聚合结果。
延迟队列解决的是数据过期问题。数据进来后不立即参与计算,而是先放进延迟队列,等达到指定延迟时间再发送回计算框架。这么做的原因很实际:某些事件的上游数据可能延迟到达,如果数据一进来就参与窗口计算,会导致窗口内明细不完整,算出来的特征值偏低。延迟队列相当于给数据加了一个等待期,确保窗口内能拿到完整的数据集。具体做法是存储 Offset 信息,系统时间与事件时间对比,超出延迟窗口期才提交 Offset,这样数据会按时间戳排好序重新进入计算。
顺序队列解决的是窗口内明细的有序性问题。队列原则是先进先出、不允许插队。每条数据进来分配一个序号(setK),用 zSet 结构管理,窗口滑动时通过 zRemRangeByRank 把头部的过期数据移除。以 Redis ZSet 为例,member 存数据明细的引用,score 存事件时间戳,每次计算窗口特征时,只需要从当前游标位置向后扫描到窗口结束位置,就能拿到完整且有序的窗口内明细。
3.2 去重计算:为什么不需要存储全部明细
PDF 里特别问了一个问题:自然窗口、固定窗口再做去重计算时,必须要存储明细数据吗?答案是如果你用的框架支持状态管理,可以用近似去重或者哈希结构来压缩存储,但如果你用的是自研框架或者简单的 Redis 方案,去重计算确实绕不开明细存储。我的做法是分维度看:如果是亿级用户、每个用户每小时的行为量在百条以内,明细存 Redis 是可接受的;但如果你要算的是 IP 级别的去重设备数,明细存储很快就把内存打爆。
常见的优化策略包括:在进入窗口计算前先按去重键做 map 阶段去重,把相同 key 的重复记录合并成一条;对 long 型 ID 用 bitmap + RoaringBitmap 做压缩存储,减少内存占用;窗口滑动时只增量更新差值部分,不重算全量。用 RoaringBitmap 做设备去重,一个亿级 ID 的集合压缩后通常只需要几十 MB 内存,相比存原始明细节省了 90% 以上。从工程实践的角度说,去重计算最怕的不是计算复杂度,而是存储方案没设计好,导致状态无限膨胀。
3.3 字段提取与转换:多数据源字段统一
特征源不止一种。用户注册信息走业务库,行为日志走埋点,设备信息走 SDK 上报,有的数据需要解析 JSON 字段,有的数据需要做 IP 转地域,有的需要枚举值映射。在 TC 框架里,这层工作在进窗口之前完成,通过一个可配置的字段转换层处理。新增数据源时不用改计算逻辑,配置好字段映射规则后自动解析。落库时统一成特征字典里的标准字段名,避免下游特征开发到处写解析代码。
提示:这里是很多特征平台后期维护成本失控的根源。字段转换规则如果不集中管理,每个特征任务各写一套 UDF,后期排查数据问题时你会发现同一个字段有六种解析方式。
4. 避坑指南:滑动窗口误差、元数据混乱与状态膨胀的踩坑记录
4.1 滑动窗口计算偏差:Spark Streaming 微批边界与事件时间错位
现象:用 Spark Streaming 做 1 小时滑动窗口特征,窗口步长 5 分钟,结果和离线 Hive 算出的特征值对不上,偏差率在 5% 到 15% 之间波动。原因:Spark Streaming 的窗口是按 Batch 对齐的,窗口的起止时间由 Batch 间隔决定,与数据本身的事件时间不完全对齐。另外,Spark Streaming 的窗口计算是左闭右开,数据的事件时间和窗口边界判断存在几秒到几十秒的延迟,这个误差在短窗口上尤其明显。解决:改用自研 TC 框架的延迟队列 + 顺序队列,数据按事件时间排序,窗口边界严格按事件时间对齐,计算完成后做抽样校验,将离线特征和在线特征在同一时间切面上做对比,误差降到 0.5% 以内。
4.2 去重状态膨胀:明细存储导致 Redis 内存打满
现象:上线「最近 24 小时去重设备数」特征之后,Redis 内存三天涨了 20GB,触发内存淘汰策略,部分特征查询开始超时。原因:特征计算用了最简单的 SET 结构存储去重 ID,每个 ID 是 32 位字符串,亿级用户量级下,这个方案的内存开销完全不可控。而且窗口滑动时只做增量添加,旧数据没有主动清理机制。解决:换成 RoaringBitmap 存储去重 ID,窗口滑动时用 zRemRangeByRank 清理过期分片,离线数据验证去重一致性之后切流上线。上线后内存占用下降约 87%,特征查询 P99 延迟从 45ms 降到了 18ms。
4.3 特征值口径不一致:离线特征和在线特征对不上
现象:同一个特征「用户近 7 天登录次数」,离线数仓算出来是 12,在线特征系统返回的是 8。业务方不知道信谁,风控策略不敢上线。原因:离线特征按自然日跑批,T+1 产出,而在线特征按滑动窗口实时计算,两者统计的时间范围和包含数据的事件时间完全不在一个切面上。解决:统一统计口径,在数据字典里给每个特征定义明确的时间窗口语义和事件时间切面,离线特征增加当日实时分区,在线特征增加 T+1 校准任务。两边数据做小时级比对,差异超过阈值时触发告警。
4.4 数据延迟导致特征值被低估
现象:用户短时间内频繁操作,但特征系统算出来的行为计数总是比实际业务量少。原因:上游埋点数据存在秒级到分钟级的延迟,数据进来直接参与窗口计算时,窗口尾部会漏掉一部分事件;窗口往前滑动之后,这部分延迟数据就不会被补算进来了。解决:在 TC 框架里给每个数据源配置延迟时间参数,常见做法是延迟 1~2 个窗口周期再参与计算,比如窗口长度 5 分钟,延迟时间设置为 1 分钟。这样既不会等太久导致特征值滞后,又能覆盖绝大多数延迟数据的到达时间。
5. 从数据字典到统一服务:特征系统的工程化落地与性能调优实战
5.1 数据字典服务:元数据统一是流批一体的前提
PDF 里那个反问很值得琢磨:如果 Hive 没有 MetaStore 会怎样?答案是没有元数据,数据就是一堆不可解析的二进制。特征系统里的数据字典,作用相当于 Hive 的 MetaStore,把所有消息队列的 topic、字段含义、数据类型、枚举取值、窗口语义统一管理起来。实时数据仓库和离线数据仓库共用一套数据字典,哪个字段被修改,消费端能感知到更新。
在落地时,可以用 Hive Metastore 作为底层元数据存储,在此基础上扩展特征字典表,记录每个特征的名称、维度、时间窗口类型、计算逻辑、关联数据源。特征开发做完了先注册到数据字典,系统自动生成特征计算任务和特征服务,不用人手动写代码。数据字典更新用版本号管理,每次更新不影响在线运行的特征实例。
5.2 统一服务层:50ms 延迟目标的架构保障
特征计算分为离线预计算和在线实时计算,统一服务层要做的是把两类特征的访问方式统一。离线特征写入 KV 存储(比如 HBase 或 Redis),在线特征实时计算完成后写入本地缓存,服务层提供统一的特征查询 API。风控引擎发起请求时,服务层并发拉取相关特征值,在 50ms 内返回结果。
性能优化上有几个关键参数:特征结果缓存时间设在 100~300ms,避免同一用户短时间内的重复计算;对热 key 做了本地缓存分层,第一层本地 Caffeine 缓存,命中率能做到 70% 以上,未命中时去查远端 KV 或触发实时计算;特征值做了协议压缩,传输层用 Protobuf 编码,减少 CPU 开销。以 58 同城的体量,在线特征服务日常 QPS 支撑到几万到几十万级别,P99 延迟控制在 50ms 以内,全靠这套分层缓存 + 并发读取设计。
提示:如果想验证特征服务性能,可以做一个压测脚本,模拟不同并发用户的行为序列,观察特征查询的 P99 延迟和缓存命中率。重点关注 P99,而不是平均延迟,平均延迟很容易被大量缓存命中拉低,掩盖真实性能问题。
5.3 特征上线与校验:从开发到全量发布的流程
特征生产不能一上来就全量上线。首先做样本验证,拿历史数据回放,把特征计算结果和离线计算结果做比对,校验窗口边界和聚合逻辑是否一致。然后做小流量灰度,把特征服务接入风控策略的测试环境,用真实流量观察特征值分布是否符合预期。异常检测方面,设定特征值上下限和波动幅度阈值,比如设备数特征突然跌到 0 或者冲高 10 倍,都要触发告警。全量上线后持续做离线在线数据比对,差异率超过 1% 就回滚并排查原因。
6. 进阶技巧:把在线特征做成模型训练可复用的数据资产
在线特征系统不止服务于实时风控策略,把在线特征做离线回放,能直接喂给模型训练做样本构造。这个思路是:在线特征每次计算的中间结果(窗口聚合值、明细数据、去重计数)全部落到日志里,入 Kafka,由离线任务消费生成特征宽表。这样一来,模型训练用的特征和在线推理用的特征完全同源同口径,不会出现训练数据分布和线上特征分布不一致的问题。
有一个细节值得关注:在做特征回放时,要处理时序穿越问题。训练样本里某条样本的特征值只能用该样本事件时间之前的数据计算,不能用之后的数据,否则模型会被「未来数据」污染,线上效果直接崩。解决办法是按照事件时间做窗口切片,把每个时间点计算出的特征快照落下来,训练时按快照取特征,保证特征和标签的时序一致性。
做在线特征和模型特征的衔接,是比较容易踩坑的环节,我自己经历过一次模型上线后特征分布偏移的问题,当时花了一周排查,最后发现是离线特征回放的窗口边界比在线计算晚了 10 分钟。从那以后,每次新特征上线我都会强制走一遍「时序一致性检查」,确认在线计算和离线回放用的是同一套时间窗口语义再发布。在线特征系统最核心的不是框架多强,而是口径、时序、元数据是否统一,这三个维度把控好,系统就不会出大问题。希望这份 58 同城的实践拆解能帮你少走弯路——做特征系统的人,都不该被那几个窗口坑第二遍。
本文还有配套的精品资源,点击获取