特征计算与转换引擎架构:批流一体、请求时计算与 One Definition Two Materializations
特征工程占实时机器学习系统 70–80% 的工程复杂度(Conduktor 估算),而其中「计算」又是复杂度最高的一环:同一份特征既要以天/小时级批量喂给离线训练,又要以秒/毫秒级流式喂给在线推理,还要在请求发生时叠加无法预计算的上下文。本文以本大模型对多源材料的体系化整合为基础,拆解特征计算与转换引擎的架构范式——批式、流式、请求时三类变换,统一特征定义如何消除训练-服务偏差,Push/Pull 物化模式的权衡,变换 SDK 与类型系统,以及业界主流实现(Uber Michelangelo / Tecton / Feast / FeatHub / RisingWave / Hopsworks / Databricks)的选型对照与落地路线。
一、核心问题:计算不一致是特征平台的第一性难题
多数 ML 团队的起点是 ad-hoc 脚本:训练用 Python 写一遍,服务用 Java 再写一遍,两套逻辑随时间漂移,产生训练-服务偏差(Training-Serving Skew)。AI Solutions Wiki 将其归纳为三种根源:代码路径不同(批用 SQL、流用 Flink,逻辑发散)、时间语义不同(批用处理时间、流用事件时间,迟到事件处理不一致)、数据源不同(批读数仓、流读事件流)。偏差不报错,只让模型精度在生产中静默衰减,极难诊断。
特征平台的架构答案不是「再加一个计算系统」,而是把计算逻辑收敛到一份声明式定义,由平台负责把同一份定义物化(Materialize)到离线存储与在线存储两条通道,即所谓的 One Definition, Two Materializations。Uber Michelangelo 的创始团队后来创立 Tecton,正是把这个思路做成托管产品;Feast 则以更轻的 registry+serving 分层切入。下面的分析围绕「计算引擎如何支撑这一范式」展开。
二、三类变换范式:批式、流式、请求时
Databricks 对特征平台给出的权威三分法是理解计算引擎架构的锁钥:
| 变换类型 | 作用对象 | 典型数据源 | 代表特征 | 计算引擎 |
|---|---|---|---|---|
| Batch(批式) | 静止数据(at rest) | 数仓 / 数据湖 / 数据库 | 商户日均交易额,每日更新 | Spark / Hive / dbt / Snowflake |
| Streaming(流式) | 事件流 | Kafka / Kinesis / PubSub / Flink | 用户近 30 分钟交易笔数,每秒更新 | Flink / Kafka Streams / ksqlDB |
| On-Demand(请求时) | 仅预测时可得的数据 | 推理请求 payload / 实时上下文 | 当前交易额是否高于用户均值 2σ | 特征服务内 Python/SQL 函数 |
三者并非互斥,而是按「新鲜度需求」与「是否可预计算」逐级叠加。CSDN《特征工程全景实战方法论》给出工业界主流选择是批流混合:离线画像用 Spark/Hive,实时行为用 Flink,通过 Feature Store 统一特征定义,确保两路输出一致。
三、批式计算引擎架构:回填、时点正确与物化
批式引擎负责历史特征回填(backfill)与训练数据生成,读 Hive/Iceberg/数据湖,写离线存储(Parquet on S3、BigQuery、Delta Lake)。其核心技术难题是Point-in-Time(PIT)查询:对每条训练样本,只能取该时刻之前「可用」的特征值,否则构成数据泄漏。Feast 的离线存储通过事件时间戳做 time-travel join 来保证 PIT 正确;Data Engineering Weekly 的 Kafka+Feast+BigQuery 模式即把 BigQuery 作为离线 materialization 的规范源,用可重复的 materialize job 同步特征。
批式物化通常通过 feast materialize-incremental 按 cron 调度(TildAlice 实测 5 万用户 3–8 分钟),或 Tecton 的托管 Spark/Rift 作业自动编排。关键架构决策是批式窗口与流式窗口对齐——同一聚合(如 7 日交易额)在批与流中必须用同一窗口语义,否则即便数值接近也会因边界不同产生偏差。
四、流式计算引擎架构:窗口、流表 Join、双写与一致性
流式是实时特征的核心。AI Solutions Wiki 的 Real-Time Feature Computation Pattern 给出标准拓扑:Events → Stream Processor → Online Store(推理)→ Offline Store(训练)。Flink 实现要点有三:
- 窗口聚合:滑动窗口(如 30 分钟滑动、1 分钟步长)用
SlidingEventTimeWindows;Flink 的MultiFeatureAggregator在同一窗口内一次算出多特征,避免重复处理。 - 流表 Join:交易事件需关联用户账龄/信用分,小维表用 broadcast state,大维表用 async I/O(
AsyncDataStream.unorderedWait,超时 1s、并发 100)避免阻塞流。 - 双写:计算结果同时写在线存储(Redis/DynamoDB)与离线存储(S3/Delta),保证训练-服务一致性。
Lambda vs Kappa:Lambda 保留批+流两套逻辑(各取所长,但一致性靠人工维护);Kappa 只用流(一套逻辑,但超长窗口效率不如批)。工业界折中是批流混合 + 统一特征定义。RisingWave 提供第四条路径:用流式数据库把特征变换写成 SQL 物化视图,增量维护、秒级更新、PG 协议直查,省去 Flink 写外部存储的中间跳。
一致性层面,Kafka Streams 在多数场景可做到 exactly-once;若在线存储用幂等 upsert(带时间戳键),至少可容忍 at-least-once + 去重。这正是 Data Engineering Weekly 强调的「compaction topic + 消息信封含 event_id/schema_version」的落地细节。
五、请求时计算(On-Demand Feature View):把变换推到 Serving API
这是 2025–2026 年增长最快的子架构。On-Demand 特征不预物化、不落在线存储,而在推理请求时由特征服务同步执行变换函数(Tecton ODFV、Feast on_demand_feature_view)。它解决的是「无法预计算」的上下文:
- 请求时才可得的数据(用户当前 GPS、本次交易额);
- 组合爆炸特征(两用户间距离,10M 用户交叉需 100 万亿值,不可预存);
- 无需重物化的新特征(在 Tecton 聚合上叠加后处理、空值填补)。
架构约束有四点(Inferensys 归纳):无状态纯函数(同输入同输出、无副作用、单 digit 毫秒内完成);混合检索(一次请求并行取预存特征 + 执行 ODFV,合并后给模型);冷启动缓解(新用户无历史,用会话上下文提取信号);延迟预算(ODFV 叠加 ~5–15ms 于基础查询 10–20ms 之上,须监控 P99 单独计)。
六、变换 SDK 与类型系统:声明式定义的工程化
特征是「一等公民」的前提是可版本化、可复用的声明式定义。Tecton 把变换抽象为可独立于 Feature View 的 @transformation 对象(模块化、UI 可发现、跨视图复用),并支持多 mode:spark_sql / snowflake_sql / pyspark / pandas / python / bigquery_sql。输出 schema 须显式含 join keys + 时间戳 + 特征值列。
类型系统是可靠性的关键:Tecton 用 tecton.types(Field("amount", Float64)、Bool、String)做请求 schema 与输出 schema 约束;ODFV 输出须非空。FeatHub(Flink 社区)则用统一 SDK 表达「所有已知特征计算逻辑」,下层可翻译为 Local / Flink / Spark 三种执行引擎——同一份定义,单机实验与分布式生产共用,是 One Definition Two Materializations 的另一种实现。
| 平台 | 计算重心 | 流式原生 | 请求时计算 | 类型系统 | 托管程度 |
|---|---|---|---|---|---|
| Feast | Registry + Serving | 需外接 | Client/On-Read 侧 | ValueType(轻) | 自管开源 |
| Tecton | 全托管 Compute | 原生一流 | Server 侧 ODFV | tecton.types(强) | SaaS/混合 |
| FeatHub | 流批一体 SDK | 原生 Flink | Feature Service 在线算 | Table Descriptor | 自管/开源 |
| RisingWave | 流式 SQL 物化视图 | 原生 | 视图即特征 | PG 类型 | 自管 |
| Databricks | Unity Catalog 集成 | Delta Live | On-Demand 三类 | DataFrame | 云托管 |
七、物化模式:Push vs Pull 的延迟-一致性权衡
Feast 文档明确采用Push 模型:数据生产者把特征值推入在线存储,读取路径只读优化,延迟最低(避免逐一拉取各生产者的网络跳数)。代价是强一致性不默认可用,需在「编排 Feast 更新」与「客户端使用」两侧显式设计。三种写模式(implicit push / write-time / on-demand read)各有取舍:写入时转换读取快但存储占用高,读取时转换省空间但延迟高。这与 ODFV 的 On-Read vs On-Write 是一体两面的同一权衡。
八、六维决策权衡表
| 决策点 | 选项 A | 选项 B | 权衡 |
|---|---|---|---|
| 计算归属 | 平台托管(Tecton) | 自管(Feast+Spark/Flink) | A 省运维但贵且有锁定;B 灵活但需基建力 |
| 流架构 | Kappa(纯流) | Lambda(批+流) | A 一套逻辑;B 长窗口更高效 |
| 实时特征 | 预物化流式 | 请求时 ODFV | 预物化低延迟高存储;ODFV 省存储加延迟 |
| 物化 | Push | Pull | Push 读优但弱一致;Pull 简单但慢 |
| 在线存储 | Redis | DynamoDB/Bigtable | Redis 低延迟小集合;后者高基数高成本 |
| 类型约束 | 强类型(tecton.types) | 弱类型(ValueType) | 强类型防漂移;弱类型快上手 |
九、30 / 60 / 90 天落地路线
- 0–30 天(止血):先建离线存储 + 批式计算,统一特征定义为声明式 registry,立即消除训练-服务偏差与重复开发;用 Feast 或既有 Spark/dbt 接入。
- 30–60 天(实时化):为欺诈/推荐/定价等低延迟场景引入流式引擎(Flink/Kafka Streams),做窗口聚合 + 双写,对齐批流窗口语义;接入在线存储(Redis)。
- 60–90 天(智能化):对无法预计算的上下文引入 ODFV,设延迟预算与 P99 监控;补齐特征监控(新鲜度、完整性、分布漂移),上线 RBAC 与血缘。
十、演进主线与反模式红旗
演进主线:ad-hoc 脚本 → 统一 registry → 批流一体 → 请求时计算 → 全托管平台(Databricks 2025 收购 Tecton,标志托管化收敛)。反模式红旗:① 两套特征代码(训练/服务各一)必漂移;② 用处理时间替代事件时间导致迟到数据错乱;③ 预物化一切组合特征(存储爆炸);④ ODFV 内调用外部 API 击穿延迟预算;⑤ 忽视特征新鲜度监控(在线特征腐坏静默拖累所有模型);⑥ 长窗口盲目走 Kappa 牺牲效率。
参考来源
- RisingWave — Real-Time Feature Engineering for ML Pipelines
- Init House — AI Data Pipelines 2026
- Data Engineering Weekly — Kafka, Feast & BigQuery 2026
- AI Solutions Wiki — ML Feature Platform
- AI Solutions Wiki — Real-Time Feature Computation Pattern
- 侯瑞哲 — 批流特征生产架构深度解析
- CSDN — 特征工程全景实战方法论
- Flink 中文社区 — FeatHub 流批一体实时特征平台
- Conduktor — Real-Time ML Pipelines
- Feast 官方文档 — Introduction / Push vs Pull:原文 · 原文
- TildAlice — Feast vs Tecton 延迟与成本基准
- AI Solutions Wiki — Feast vs Tecton 对比
- ML Journey — Feast vs Tecton vs Hopsworks
- Tecton 文档 — On-Demand Feature View
- Tecton 文档 — Building On-Demand Features
- Neel Mishra — Tecton Enterprise Feature Platform
- Tecton 文档 — Feature Design Patterns / Transformations:原文 · 原文
- Databricks — What is a Feature Platform
- Tecton Blog — How Do Real-Time Features Work:原文 · 原文
- CSDN — Feast 特征转换 On-Demand 实时计算
- Inferensys — On-Demand Features