你点击「立即申请」的那 100ms,风控系统在做什么

你打开某个 App,看到一条"最高借 5 万,秒到账"的信贷广告,点了"立即申请"。

从你的手指离开屏幕,到页面显示"审核通过,额度 2 万元"——这中间大概不到一秒。但在这一秒里,系统完成了一件相当复杂的事:它从几十个数据源里拉取了你过去数年的行为数据,实时计算了上百个特征,把这些特征输入到一个训练了数亿样本的风控模型,给出了一个信用评分,再根据评分决定你能不能借、能借多少。

这整条链路,叫做在线风险决策系统。模型平台工程师建设的,正是支撑这条链路运转的基础设施。

这篇文章把这条链路从头到尾走一遍:数据从哪来、特征怎么算、模型怎么训练、推理怎么做到足够快。跟着走完,你会明白 Hudi、HBase、Redis、Flink、Spark 各自在这条链路的哪个位置、解决了什么具体问题。

一、全局视图:系统长什么样

先看整体。一个典型的金融 ML 系统,从数据到决策,大概分四个层:

数据源层 MySQL/OLTP Kafka 消息队列 Hive 数仓 外部数据/日志 特征工程层 批式计算 Spark / Hive SQL Offline Store Hudi on HDFS 流式计算 Flink Online Store Redis / HBase 模型训练层 读 Offline Store → 生成样本 PyTorch / TensorFlow / 自研框架 模型推理层 读 Online Store → 实时推理 Triton / TorchServe / 自研推理服务 模型文件 → 上线部署 风险评分 → 业务决策
金融 ML 系统的四层结构:数据源 → 特征工程 → 训练/推理 → 决策输出

注意这张图里有两条并行的路径:左边是批式路径(离线),数据通过 Spark 处理后存入 Hudi,用于模型训练;右边是流式路径(在线),数据通过 Flink 实时计算后存入 Redis/HBase,用于线上推理。两条路径算的是同一批特征,但服务于不同用途。

为什么要这样分?接下来一步步解释。

二、第一步:原始数据在哪里

用户申请信贷,这个事件发生时,系统里同时存在几类数据:

  • 结构化数据库(MySQL/OLTP):用户的基础信息——注册时间、认证状态、历史还款记录。这些数据是"静态"的,变化频率低,一天更新几次。
  • 消息队列(Kafka):用户的实时行为事件——刚刚打开了 App、浏览了贷款页面、点击了申请按钮。这些数据是"流式"的,毫秒级产生。
  • 数据仓库(Hive):已经过 ETL 处理的历史数据——过去 90 天的消费记录、历史贷款金额、逾期次数。这些数据是"批量"的,通常是 T-1 的,昨天的数据今天才能用。

这三类数据来自不同系统,格式不同,延迟不同,可靠性保证也不同。特征工程的第一件事,就是把它们整理成模型能用的格式。

三、第二步:批式特征计算(Spark + Hudi)

每天凌晨,有一批 Spark 任务在跑。它从 Hive 里读取昨天的用户行为数据,计算出一批特征:

  • 用户过去 7 天的贷款申请次数
  • 用户过去 30 天的平均消费金额
  • 用户过去 90 天的最大单笔交易
  • 用户历史逾期天数总和

这些特征计算出来之后,写入 Hudi(Hadoop Upsert Delete Incremental)。

为什么不直接写回 Hive?因为 Hudi 解决了两个 Hive 解决不了的问题:

第一,支持行级更新(Upsert)。 Hive 是不可变的,已经写入的分区数据不能修改。但用户的特征每天都要更新,Hudi 支持按主键 Upsert——相同用户 ID 的记录会被覆盖更新,而不是追加新行。

第二,支持时间旅行查询(Time Travel)。 Hudi 保留了数据的历史版本,你可以查询"2024 年 1 月 15 日这个用户的特征值是多少"。这个能力对模型训练至关重要——训练样本需要的是事件发生时刻的特征,而不是当前最新的特征(否则就是数据泄漏)。

批式特征的特点:准确、完整、但有延迟。凌晨 3 点算完,覆盖到用户上午申请时已经是"昨天"的数据了。这对风控来说还不够——欺诈分子可以在几分钟内发起大量申请,你需要知道他刚才做了什么。

四、第三步:流式特征计算(Flink + Redis)

除了批式特征,还有一批 Flink 任务在实时消费 Kafka 的事件流。同样是"用户过去 1 小时内的申请次数"这个特征,Flink 用滑动窗口算子实时维护,每来一个新事件就更新一次。

算完的结果写到 Redis(或 HBase)。

为什么用 Redis?因为线上推理对延迟要求极高——100ms 内给出评分,其中特征读取能用的时间不超过 10ms。Redis 是内存数据库,单次读取耗时在 1ms 以内,满足要求。

什么时候用 HBase?当特征是"宽表"形式——一个用户有几百个特征字段,总大小超过 Redis 单条记录的限制时,或者数据量大到 Redis 内存装不下时,用 HBase(基于 HDFS 的 KV 存储,容量大但延迟稍高,通常在 5-10ms)。

流式特征的特点:实时、但计算逻辑相对简单。复杂的统计特征(比如用户历史行为的分布统计)在流式计算里很难精确实现,通常只做滑动窗口聚合。

现在你明白为什么需要批式 + 流式两条路径了:批式负责精确的历史统计特征,流式负责实时的行为特征,两者互补。

五、第四步:模型训练(从 Offline Store 生成样本)

模型每隔一段时间需要用新数据重新训练。训练样本的生成过程大概是:

  1. 找出一段时间内所有的信贷申请事件(正样本 = 后来逾期的,负样本 = 正常还款的)
  2. 对每条申请事件,找出该用户在申请发生时刻的特征值
  3. 把事件标签(是否逾期)和特征值拼在一起,形成训练样本

第 2 步是关键,也是容易出错的地方。

一个申请事件发生在 2024 年 1 月 15 日,你需要的是用户在那一天的特征值,而不是今天的特征值。如果用今天最新的特征(比如用户最近三个月的消费记录),你的模型就"看到了未来"——训练时学到的是未来的信息,线上推理时没有这些信息,模型效果会虚高。这就是前面提到的**数据泄漏(Data Leakage)**。

Hudi 的时间旅行查询在这里发挥作用:给定申请事件的时间戳,查询那个时刻的特征快照。这个操作叫做 Point-in-Time Correct Join,是模型平台工程里最重要的正确性保证之一。

样本生成完毕后,送入训练框架(PyTorch/TensorFlow/自研)。训练完成的模型文件(通常是 SavedModel 格式或 TorchScript 格式)存储到模型仓库,等待上线。

六、第五步:线上推理(100ms 内完成)

用户点击"申请"按钮,后端接到请求,触发风险评分。这个过程大概是:

👤 申请请求 特征读取 Redis: 实时特征 HBase: 历史特征 ≈ 5-10ms 特征拼装 归一化/拼向量 模型推理 Triton/TorchServe GPU/CPU 计算 ≈ 20-50ms 风险评分 0.0 - 1.0 总延迟目标:100ms 以内
一次线上推理的完整链路,从用户请求到风险评分,100ms 内完成

每个步骤的延迟预算:

  • 特征读取(Redis + HBase):5-10ms
  • 特征拼装和预处理(归一化、缺失值填充):1-2ms
  • 模型推理(神经网络前向传播):20-50ms(取决于模型大小和硬件)
  • 网络和其他开销:剩余预算

这就是为什么特征读取必须用 Redis 而不是数据库——数据库的查询延迟通常在 5-50ms,高峰期可能更高,而且无法承受风控系统的高并发请求量。

七、最重要的约束:训练和推理必须用同一套特征

到现在你已经理解了整个系统。但有一个约束贯穿全程,是这个系统最核心的正确性保证:训练用的特征和线上推理用的特征,必须是同一套计算逻辑产生的

这听起来是废话,但实际上极难保证。

训练侧用 Spark SQL 算"用户过去 7 天的点击次数",推理侧用 Flink 算同一个特征。两套代码,分别维护,随着版本迭代,细节开始不同:Spark 按自然天统计,Flink 按滚动窗口统计;Spark 统计完整的历史数据,Flink 只统计当天已到达的事件;Spark 有去重逻辑,Flink 没有……

结果:模型训练时学到的是"点击次数 = 50 时,风险高",但线上给模型的是 Flink 算出的 65。模型的决策边界完全错位,但整个系统不会报任何错误,只是评分慢慢变得不准。

这个问题叫做 Training-Serving Skew,是 ML 系统最难排查的 bug 类型之一,也是你作为模型平台工程师最需要关注的事情。

解决它的根本思路是:不维护两套逻辑,让批式和流式从同一份特征定义生成。具体实现有几种方向:Flink 批流一体(同一份 SQL,切换执行模式)、统一特征定义语言(用 YAML 描述特征逻辑,系统自动翻译成 Spark 和 Flink 代码)、或者至少做到特征计算逻辑共享代码库而非各自实现。

八、你将要建设什么

现在你能把模型平台的各个工作职责对应到系统的具体位置了:

  • "训练框架开发":第五步,训练侧的基础设施——样本生成流程、Point-in-Time Join、分布式训练任务管理
  • "在线模型推理框架":第六步,推理侧——推理服务的性能优化、模型热更新、请求调度
  • "特征平台建设":第二到四步,整个特征工程层——批式特征 Pipeline、Flink 实时特征、Offline/Online Store 的管理和一致性保证
  • "批、流式特征生产架构":第三步和第四步,Spark 批式 + Flink 流式的完整架构
  • "数据隔离平台":在多个业务线共用一个特征平台时,保证 A 业务的数据不会影响 B 业务的计算结果

你的核心工作不是训练模型,而是让算法工程师能高效地训练和部署模型。特征平台越好用、训练速度越快、推理延迟越低,算法工程师就能更快地迭代,模型效果就越好。这是一个间接但非常关键的角色。

九、有了这张地图,从哪里开始

现在你有了整体认知,具体准备可以有针对性:

动手跑一遍 Hudi。 在本地搭一个最简单的 Hudi demo:写入几条数据、做一次 Upsert、然后用 time-travel query 查历史版本。这个操作本地一两小时能跑通,但会让你对 Hudi 的工作方式有直觉。

看懂一个 Flink 滑动窗口 demo。 找一个 Flink WordCount 的变体,改成"统计最近 1 小时内每个用户的点击次数"。跑通它,理解 Event Time、Watermark、State 是怎么配合的。

用 PyTorch 本地跑一个推理。 找一个预训练的简单模型(比如分类模型),用 Python 加载后做推理,记录延迟。然后换成 TorchServe 部署后再测一次。感受一下推理框架解决了什么问题。

理解业务场景。 搜索一下"美团金融风控"、"消费金融 ML 系统"相关的技术文章(美团技术博客有不少)。理解业务方向,第一次和产品、算法同事讨论需求时才能接上话。

不需要全部精通,但每个方向有一点亲身操作的经验,比看十篇文章更有用。

十、总结

用一张地图总结整个系统:

  • 数据源:MySQL(静态)、Kafka(实时)、Hive(历史批量)
  • 批式路径:Spark → Hudi(Offline Store)→ 模型训练
  • 流式路径:Flink → Redis/HBase(Online Store)→ 线上推理
  • Point-in-Time Join:训练时用事件发生时刻的特征,防止数据泄漏
  • Training-Serving Skew:批流两套逻辑的不一致,是最难排查的 bug
  • 推理延迟分解:特征读取(5-10ms)+ 模型推理(20-50ms)= 100ms 目标

你建设的不是模型本身,而是让模型能被快速训练、准确评估、高效服务的基础设施。做好这件事,算法同事才能专注于算法。