你打开某个 App,看到一条"最高借 5 万,秒到账"的信贷广告,点了"立即申请"。
从你的手指离开屏幕,到页面显示"审核通过,额度 2 万元"——这中间大概不到一秒。但在这一秒里,系统完成了一件相当复杂的事:它从几十个数据源里拉取了你过去数年的行为数据,实时计算了上百个特征,把这些特征输入到一个训练了数亿样本的风控模型,给出了一个信用评分,再根据评分决定你能不能借、能借多少。
这整条链路,叫做在线风险决策系统。模型平台工程师建设的,正是支撑这条链路运转的基础设施。
这篇文章把这条链路从头到尾走一遍:数据从哪来、特征怎么算、模型怎么训练、推理怎么做到足够快。跟着走完,你会明白 Hudi、HBase、Redis、Flink、Spark 各自在这条链路的哪个位置、解决了什么具体问题。
一、全局视图:系统长什么样
先看整体。一个典型的金融 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 生成样本)
模型每隔一段时间需要用新数据重新训练。训练样本的生成过程大概是:
- 找出一段时间内所有的信贷申请事件(正样本 = 后来逾期的,负样本 = 正常还款的)
- 对每条申请事件,找出该用户在申请发生时刻的特征值
- 把事件标签(是否逾期)和特征值拼在一起,形成训练样本
第 2 步是关键,也是容易出错的地方。
一个申请事件发生在 2024 年 1 月 15 日,你需要的是用户在那一天的特征值,而不是今天的特征值。如果用今天最新的特征(比如用户最近三个月的消费记录),你的模型就"看到了未来"——训练时学到的是未来的信息,线上推理时没有这些信息,模型效果会虚高。这就是前面提到的**数据泄漏(Data Leakage)**。
Hudi 的时间旅行查询在这里发挥作用:给定申请事件的时间戳,查询那个时刻的特征快照。这个操作叫做 Point-in-Time Correct Join,是模型平台工程里最重要的正确性保证之一。
样本生成完毕后,送入训练框架(PyTorch/TensorFlow/自研)。训练完成的模型文件(通常是 SavedModel 格式或 TorchScript 格式)存储到模型仓库,等待上线。
六、第五步:线上推理(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 目标
你建设的不是模型本身,而是让模型能被快速训练、准确评估、高效服务的基础设施。做好这件事,算法同事才能专注于算法。