机器学习管道
在我最喜欢的《辛普森一家》(The Simpsons)剧集中,当霍默·辛普森听说培根、火腿和猪排都来自同一种动物时,他简直不敢相信:「没错,丽莎,一只神奇而美妙的动物。」当我向 ChatGPT 4.1 询问 ML 管道(ML pipeline)的定义时,我也有同样的反应。它告诉我,ML 管道执行数据收集、特征工程、模型训练、模型评估、模型部署、模型监控、推理和维护。「没错,GPT,一条神奇而美妙的单体 ML 管道,」我想。它甚至声称它的 ML 管道是模块化的!
难怪当我问 10 位不同的数据科学家什么是 ML 管道时,我通常会得到 10 种不同的答案。对于它的输入和输出是什么,人们没有达成共识。如果一位开发者告诉你,他们用 ML 管道构建了他们的 AI 系统,你能从中获得什么信息?在我看来,ML pipeline(ML 管道)这个术语,按照目前的使用方式,在讨论构建 AI 系统时可能「被认为是有害的」(considered harmful)。1 在本书中,我们力求更加严谨。我们用构建 AI 系统时所使用的具体管道来描述 AI 系统。我们保留 ML pipeline 这个术语的使用,用它来描述 AI 系统中的任何单个管道或一组管道。
管道(pipeline)是一个具有明确定义的输入和输出的计算机程序(也就是说,它有一个定义良好的接口),并且按计划运行或持续运行。ML 管道是任何输出用于 AI 系统的 ML 工件(ML artifact)的管道。我们以 ML 管道创建或修改的 ML 工件来命名具体的 ML 管道。创建 ML 工件的 ML 管道包括:输出特征的特征管道(feature pipeline)、输出嵌入的向量嵌入管道(vector-embedding pipeline)、输出已训练模型的训练管道(training pipeline),以及输出预测的推理管道(inference pipeline)。修改 ML 工件的 ML 管道包括:将模型从未验证状态转变为已验证状态的模型验证管道(model validation pipeline),以及将模型部署到生产环境的模型部署管道(model deployment pipeline)。在本章中,我们介绍许多不同类型的 ML 管道,但我们将重点深入探讨对构建 AI 系统最重要的 ML 管道——特征管道、训练管道和推理管道。三条管道与真相。
用 ML 管道构建 AI 系统
在探讨如何开发 ML 管道之前,我们先来看一个构建 AI 系统的开发流程。AI 系统是软件系统,软件工程方法论有助于指导你构建软件系统。第一代面向 ML 的软件开发流程,例如微软(Microsoft)的团队数据科学流程,主要集中在数据收集和建模上,但没有解决如何构建 AI 系统的问题。因此,这些流程很快被 MLOps(机器学习运维)所取代,MLOps 专注于自动化、版本控制以及开发人员与运维人员之间的协作来构建 AI 系统。
最小可行预测服务
我们在此介绍一种极简的 MLOps 开发方法论,其核心是尽可能快地得到一个最小可行的 AI 系统,即最小可行预测服务(minimal viable prediction service,MVPS)。我在 KTH 的构建 AI 系统课程中遵循了这一 MVPS 流程,它使学生们能够在最多几天内就得到一个可运行的 AI 系统(该系统使用新颖的数据源来解决新颖的预测问题)。
注
ML 工件包括模型、特征、训练数据、向量索引、模型部署以及预测/上下文日志。ML 工件是由 ML 管道产生的有状态对象,由你的 ML 基础设施服务管理。大多数 ML 工件是不可变的,但特征数据、向量索引和模型部署除外,它们可以在原地更新。
如图 图 2-1 所示的 MVPS 开发流程,从识别以下内容开始:
- 你想解决的预测问题
- 你想改进的 KPI(关键绩效指标)指标
- 你可以使用的数据源
一旦你确定了构成 AI 系统的这三根支柱,你就需要将你的预测问题映射到一个 ML 代理指标(ML proxy metric)上——即你将在 AI 系统中优化的目标。这通常是最具挑战性的一步。ML 代理指标还应该与 KPI 正相关。
描述最小可行预测服务迭代开发过程的示意图:从识别预测问题并将其映射到 ML 代理指标开始,随后实现特征管道、训练管道和推理管道,最后以用户界面或系统集成收尾。

接下来是实施阶段,你通常从左到右推进,但任何时候如果需要重新定义你的预测问题、KPI 或数据源,你都可以折返。实施步骤如下:
- 开发一个最小的特征管道,既能回填(backfill)历史数据,又能将增量生产数据写入你的特征存储(feature store)。
- 如果你需要自定义模型,则开发一个最小的训练管道(如果你使用预训练模型(如大语言模型(large language model,LLM)),则跳过此步骤)。
- 开发一个推理管道,用你的模型进行预测。它可以是一个批处理程序、一个在线推理程序、一个 LLM 应用或一个智能体(agent)。
- 开发一个 UI 或仪表板,让利益相关者可以试用你的 MVPS,并让你能够迭代改进它。
让我们从头开始,用一个电商商店的例子说明:你想预测用户感兴趣的商品或内容。对于电商商店中的商品推荐,KPI 可以是转化率的提升,以用户将商品加入购物车来衡量。对于内容,一个可衡量的业务 KPI 可以是用户参与度的最大化,以用户在服务上花费的时间来衡量。作为数据科学家或 ML 工程师,你的目标是接受预测问题和业务 KPI,并将其转化为一个优化某个 ML 指标(或目标)的 AI 系统。ML 指标可能直接匹配业务 KPI,例如用户将商品加入购物车的概率;或者 ML 指标可能是业务 KPI 的代理指标,例如用户将会与推荐内容互动的预期时间(这是在平台上提升用户参与度的代理指标)。
一旦你有了预测问题、KPI 和 ML 目标,你就需要思考如何基于可用的数据,创建带有对目标具有预测能力的特征的训练数据。你应该从枚举并获取为 AI 系统提供数据的数据源的访问权限开始。然后你需要理解数据,这样才能有效地从数据中创建特征。探索性数据分析(exploratory data analysis,EDA)是你通常采取的第一步,以了解你的数据、其质量,以及任何特征与目标变量之间是否存在依赖关系。如果你还不太熟悉该领域,EDA 通常有助于培养对数据的领域知识。它可以帮助你识别哪些变量可以或应该用于模型或被创建出来,以及它们对模型的预测能力。你可以通过使用 LLM 驱动的助手(如 Hopsworks Brewer)检查数据及其分布来开始 EDA,或者将数据摄取到特征存储中,由特征存储在摄取时计算数据统计信息。如果需要,你可以在笔记本中通过可视化分析和统计方法进行更详细的 EDA。
下一个(不可避免的)步骤是确定你将用于构建 FTI 管道(即特征、训练和推理管道)的不同技术(见 图 2-2)。我们建议为此使用看板(kanban board)。看板是一种可视化工具,用于跟踪工作在 MVPS 流程中的流转,包含不同阶段的列和单个任务的卡片。Atlassian Jira 和 GitHub Projects 是开发者广泛使用的看板示例。
看板示意图,说明从数据源到 AI 应用的流程,突出显示特征管道、训练管道和推理管道及其相关技术和编排器。

在开始实现 AI 系统之前,填写 MVPS 看板是一个很好的活动,可以让你对所构建的 AI 系统有一个整体概览。你应该把看板的标题设为你的 AI 系统所解决的预测问题的名称,然后填写数据源、将消费预测的 AI 应用,以及你打算用来实现 FTI 管道的技术。你还可以用非功能性需求来标注不同的看板泳道,例如特征管道的体量、速率和新鲜度要求,或在线推理管道响应时间的服务级别目标(service-level objective,SLO)。在起草了系统架构之后,你就可以开始编写代码了。你之后可能会更改所选的技术和非功能性需求,但对前进方向有一个愿景是很好的实践。
此时,你已经了解了数据和所需的特征,现在你需要从数据源中提取目标观测值(或标签)和特征。这涉及从数据源构建特征管道。特征管道的输出将是存储在特征存储中的特征和观测值/标签。如果你已有特征存储,并且幸运的是它已经包含你需要的目标(们)和/或特征,你就可以跳过实现特征管道这一步。
从特征存储中,你可以创建训练数据,然后实现训练管道来训练模型,并将模型保存到模型注册表(model registry)。最后,你实现一个推理管道,使用你的模型和新的特征数据进行预测,并添加 UI 或仪表板来创建你的 MVPS。这个 MVPS 开发流程是迭代的,你逐步改进 FTI 管道。你添加测试、验证和自动化。你之后可以添加不同的环境:开发、预发布和生产。
为 ML 管道编写模块化代码
一个成功的 AI 系统需要随着时间的推移进行更新和维护。这意味着你需要对源代码进行各种更改,例如:
- 计算的特征集或其计算所依据的数据
- 如何训练模型(其模型架构或超参数),以提高其性能或减少偏差
- 对于批量 ML 系统,更频繁(或更不频繁)地进行预测,或更改保存预测结果的目标位置(sink)
- 对于在线 ML 系统,请求延迟或特征新鲜度要求的变化
- 对于 LLM 应用和智能体,上下文工程、工具或 LLM 版本的变化
在系统架构层面,我们可以将 AI 系统模块化为我们的三个(或更多)管道——特征管道、训练管道和推理管道。这种模块化程度使你能够独立开发每个管道——只要你不破坏每个管道的数据契约(data contract)。每个管道的数据契约包括其输入/输出模式(schema)以及任何非功能性需求,例如特征管道的数据验证规则、训练管道的模型性能或偏差,或在线推理管道的 SLO。
然而,在每个 ML 管道内部,你还需要编写遵循软件工程最佳实践的模块化代码。你的源代码应该经过测试且易于维护,并且应该遵循 DRY(「不要重复自己」)原则。如果你的 ML 管道的源代码是一堆意大利面条式笔记本,将很难构建可靠的 ML 管道。你将如何测试笔记本中的代码,以确保所做的任何更改在部署到生产环境之前都能正确工作?你将如何让新开发人员上手处理代码库?
我们建议你在用 Python 编写 ML 管道时采用的方法,是将源代码重构为函数或类。你将 ML 管道中的步骤分解为一组函数,这些函数组合在一起实现 ML 管道程序。每个函数应该封装一项可管理且相关的工作,并且函数可以在代码库的不同部分重用。你将函数的实现(及其所有复杂性)隐藏在接口后面。在 Python 中,函数的接口是函数的签名(signature)——其名称、参数和返回类型。
笔记本作为 ML 管道?
最佳实践是将特征函数存储在 Python 模块中(而不是笔记本中),这样它们可以被独立地进行单元测试,并在不同的 ML 管道中重用。然而,ML 管道程序仍然可以是导入并使用特征函数的笔记本。如果你想将 ML 管道作为笔记本运行,你需要使用支持将笔记本作为作业调度的平台(例如 Hopsworks 上的 Jupyter Notebooks)。我们不推荐使用 Google Colaboratory(Colab)笔记本,因为它们与 Git 配合不佳。没有 Git 支持,很难将 GitHub 仓库中的 Python 模块文件导入到你的 Colab 笔记本中。
我们先看一些示例特征工程代码,我们想对它进行重构,使其更易于测试和维护。在下面的特征管道代码中,有一个 compute_features 函数对 Pandas DataFrame 执行数据变换。这是一个在 Pandas 中进行非模块化特征工程的示例:
import pandas as pd
def compute_features(df: pd.DataFrame) -> pd.DataFrame:
if config["region"] == "UK":
df["holidays"] = is_uk_holiday (df["year"], df["week"])
else:
df["holidays"] = is_holiday (df["year"], df ["week"])
df["avg_3wk_spend"] = df["spend"].rolling (3).mean()
df["acquisition_cost"] = df["spend"]/df["signups"]
df["spend_shift_3weeks"] = df["spend"].shift(3)
df["special_feature1"] = compute_bespoke_feature(df)
return df
df = pd.read_parquet("my_table.parquet")
df = compute_features(df)
这段代码不是模块化的,因为一个函数计算了五个特征(holidays、avg_3wk_spend、acquisition_cost、spend_shift_3weeks 和 special_feature1)。为每个单独的特征编写独立的测试很困难,没有针对每个特征的专门文档,调试也需要理解整个 compute_features 函数。
解决这些问题的方法是将这段代码重构为特征函数(feature function),即更新包含特征的 DataFrame 的函数。这个想法最初来自 Apache Hamilton。对于每个计算的特征,你定义一个特征函数。你可以通过按正确顺序应用特征函数,将特征创建为 DataFrame(Pandas、PySpark 或 Polars)中的列。例如,在这里,我们将 acquisition_cost 列计算为 spend 除以注册我们服务的用户数(signups):
df['acquisition_cost'] = df['spend'] / df['signups']
我们将用于计算 acquisition_cost 的逻辑重构为一个函数,如下所示:
def acquisition_cost(spend: pd.Series, signups: pd.Series) -> pd.Series:
"""Acquisition cost per user is total spend divided by number of signups."""
return spend / signups
我们也为其他四个特征编写函数。乍一看,这增加了我们必须编写的代码行数。然而,我们现在有了一个带有文档的函数,它可以在同一个程序内或由不同的程序重用。我们现在可以为 acquisition_cost 编写单元测试,如下所示:
@pytest.fixture
def get_spends(self) -> pd.DataFrame:
return pd.DataFrame([[20, 40], [5, 4], [4, 10],
columns=["spends", "signups", "acquisition_cost"])
def test_spend_per_signup (get_spends : Callable):
df=get_spends()
df["res"] = acquisition_cost(df["spends"], df["signups"])
pd.testing.assert_series_equal(df["res"], df["acquisition_cost"])
这个单元测试为 acquisition_cost 的计算方式强制执行了一个契约。如果任何人更改了 acquisition_cost 的计算方式,单元测试将失败,表明对于使用 acquisition_cost 特征的下游客户端,其契约已被破坏。当然,你仍然可以更新 acquisition_cost 的特征逻辑,但这通常应该通过创建特征的新版本来完成,而新版本需要新的单元测试。我们将在第5章中介绍特征版本化。
在这个例子中,我们的函数是特征管道中对 DataFrame 的数据变换。特征管道如何将最终的 DataFrame 保存到特征存储?特征存储通常提供 DataFrame API(Pandas、Apache Spark、Polars),用于将 DataFrame 摄取到特征组(feature group)中,特征组是特征存储保存数据的表。我们编写模块化特征工程的方法,是使用特征函数构建一个包含特征数据的 DataFrame(特征化 DataFrame,featurized DataFrame)(见 图 2-3)。
示意图:特征管道从新数据和回填数据创建特征化 DataFrame,写入特征组,之后训练管道和推理管道通过特征存储中的特征查询引擎访问它。

每个特征化 DataFrame 作为一次提交(commit)(追加/更新/删除)写入特征存储中的特征组。特征组存储随时间创建的可变特征集。训练管道和推理管道之后可以使用特征查询服务,从一个或多个特征组读取一致的特征数据快照,分别用于训练模型或进行预测。
在本书中,我们将采用特征函数的方法来模块化用于数据变换的 Python 代码。尽管我们之前的例子涵盖的是特征管道,但我们也将在训练管道和推理管道中遵循同样的编码实践:将数据变换封装在函数中。在下一节中,我们将看到,根据你创建的特征类型——可复用特征、模型专属特征或实时特征——某些数据变换仍然需要在训练管道和推理管道中执行。
ML 管道中数据变换的分类法
ML 管道由一系列数据变换组成。从数据源到特征,再到模型和预测,数据被不断地从一种格式变换为另一种格式,直到最终预测被客户端消费。然而,并非 ML 管道中的所有数据变换都是相同的。首先,特征存储存储的特征数据可以在许多模型中复用。这意味着,将特征数据写入特征存储的特征管道,应该执行创建可复用特征的数据变换。
然而,有些数据变换产生的特征无法跨模型复用。例如,许多 ML 框架要求你在将字符串用作输入之前,先将其变换为数值表示。这种变换被称为对类别变量(categorical variable)进行编码(encoding),它以模型训练数据集中找到的类别集为参数。如果你在两个不同的训练数据集上训练两个模型,每个数据集的类别集不同,它们对字符串的编码就会不同。因此,这种数据变换是特定于模型及其训练数据集的。类似地,对于数值变量,有些数据变换以模型的训练数据集为参数,因此无法跨模型复用。你可以使用从训练数据中的值计算出的均值/最小值/最大值/标准差来归一化或缩放数值。有些模型需要归一化的数值变量,例如梯度下降模型(深度学习),而其他模型(如决策树)则不会从归一化中受益。
另一种在特征管道之外执行的数据变换,是在实时 ML 系统中执行的实时数据变换。特征管道预计算特征,但在线模型可能需要对预测请求的参数进行数据变换。这些按需变换在在线推理管道中执行,例如在 Python 用户自定义函数(user-defined function,UDF)中。
为了应对这两个挑战,我们现在引入一个使用特征存储的 ML 管道中数据变换的分类法。该分类法将数据变换分为三类:模型相关、模型无关和按需变换。这种分类有助于告知你应该在哪个(些)ML 管道中实现该数据变换。但在介绍分类法之前,我们先介绍特征类型。
特征类型与模型相关变换
编程语言中变量的数据类型(data type)定义了该变量上合法操作的集合——非法操作将导致错误,无论是在编译时(Java 和 Rust)还是在运行时(Python)。特征类型(feature type)是数据类型的扩展,有助于理解 ML 中变量上合法变换的集合。例如,我们可以编码类别变量,但不能编码数值特征。类似地,我们可以对输入到 LLM 的字符串(类别)进行分词(tokenize),但不能对数值特征进行分词。我们可以对数值变量进行归一化、标准化或缩放,但不能对类别变量进行这些操作。
在 图 2-4 中,我们将特征类型定义为类别变量(字符串、枚举、布尔值)、数值变量(整数、浮点数、双精度浮点数)和数组(列表、向量嵌入)。在 ML 文献中,数组通常不被描述为一种特征类型。然而,它们现在在 AI 系统中无处不在,尤其是作为向量嵌入。向量嵌入(vector embedding)是浮点数或整数的定长数组,存储某些高维数据的压缩表示。列表和向量嵌入现在作为数据类型在特征存储中得到广泛支持——并且它们具有定义良好的合法变换集合。例如,取列表中最新的三个条目是对列表的合法操作,索引/查询向量嵌入也是如此。
将 ML 中的特征类型分为三类的示意图:类别型、数值型和数组型,并进一步细分出定序、定类、嵌入、定距和定比等子类。

特征类型缺乏编程语言支持;相反,它们在 ML 框架和库中得到支持。例如,在 Python 中,你可以使用 Scikit-Learn、TensorFlow、XGBoost 或 PyTorch 等 ML 框架,每个框架都有自己针对其特征类型的编码/缩放/归一化/最小-最大缩放变换实现。
这些变换是 ML 特有的。它们使特征数据与特定的 ML 框架兼容,或提高模型性能,例如改善基于梯度下降的 ML 算法收敛性的归一化。如前所述,这些变换无法在其他模型中复用,因此我们称这些变换为模型相关变换(model-dependent transformation,MDT)。这些变换依赖于模型和/或其训练数据。你不应该在特征管道中、在特征存储之前执行这些变换。相反,你应该应用两次 MDT:第一次在训练管道中创建训练数据时,第二次在推理管道中。由于训练管道和推理管道是不同的程序,你需要确保 MDT 在训练管道和推理管道中的实现之间没有偏差(skew)。如果存在偏差,你的模型可能表现不佳,而且很难找出性能不佳的原因。
MDT 的另一个问题是,变换后的特征数据不利于 EDA。例如,如果你对年收入变量进行归一化,数据就变得难以分析:数据科学家理解和可视化 74,580 美元的收入,比理解其归一化值 0.541 更容易。在特征存储中存储归一化/缩放/编码的特征数据也有问题。例如,如果你有一个存储归一化新年收入数据的特征组(表),每次在该表中添加/删除/更新行时,你都不得不重新计算所有现有的年收入特征数据,因为新数据会改变现有行的均值/标准差。这使得即使对特征组进行非常小的写入也代价高昂(这被称为写放大,write amplification)。
用模型无关变换实现可复用特征
数据工程师通常不太熟悉上一节介绍的 MDT,因为它们是 ML 特有的。数据工程师熟悉并在特征工程中广泛使用的数据变换类型是(窗口)聚合(例如某个数值变量的最大值/最小值)、窗口计数(例如每天的点击次数),以及创建最近消费时间、消费频率、消费金额(recency, frequency, monetary value,RFM)特征的任何变换。这些变换创建的特征可以在许多模型中复用,被称为模型无关变换(model-independent transformation,MIT)。MIT 在批量或流式特征管道中计算一次,它们产生的可复用特征数据存储在特征存储中,供之后的下游训练管道和推理管道使用。
用按需变换实现实时特征
如果我有一个实时 ML 系统,而计算特征所需的数据只作为预测请求的一部分可用,该怎么办?在这种情况下,我必须在在线推理管道中计算该特征,这被称为按需变换(on-demand transformation,ODT)。通常,预测请求及其参数会被记录下来以供后续使用。例如,你可能想复用相同的输入数据,为特征存储创建可复用特征数据。或者你可以将该历史数据用作 MDT 的输入。我们将在第7章中看到如何将 ODT 实现为用户自定义函数(UDF)。在线推理管道中使用的同一个 UDF 可以在特征管道中复用,以从历史数据创建可复用特征。我们的方法将防止偏差——在线推理管道中的数据变换与特征管道中的数据之间应该没有差异。
ML 变换分类法与 ML 管道
既然我们已经介绍了三种不同类型的特征,以及创建它们的三种不同数据变换(模型无关、模型相关和按需),我们就可以提出 ML 数据变换的分类法了(见 图 2-5)。我们的分类法包括:
- 产生可复用特征并存储在特征存储中的模型无关变换
- 产生特定于单个模型的特征的模型相关变换
- 需要在在线推理管道中使用请求时数据计算的按需变换,但也可以在特征管道中对历史数据计算
ML 数据变换分类法示意图,详细说明用于特征创建的模型无关、模型相关和按需变换。

在 图 2-6 中,我们可以看到分类法中的不同数据变换如何映射到我们的 FTI 管道。
机器学习管道中数据变换过程的示意图,包括特征管道、训练管道和推理管道,突出显示模型无关、模型相关和按需变换。

注意,MIT 只在特征管道中执行。然而,MDT 在训练管道和推理管道中都会执行。按需变换也在两个不同的管道中执行——在线推理管道和特征管道。批量推理管道不支持 ODT,因为它们没有请求时参数——它们的预计算特征在特征管道中计算,任何推理时变换都是 MDT。每当同一个数据变换在不同的管道中执行时,你都需要确保不同实现之间没有偏差。最后一点需要注意的是,MDT 也可以应用于在线推理管道中的请求参数,但它们与 ODT 的不同之处在于,它们不能应用于特征管道。因此,有些实时特征可以是模型无关特征,而另一些则是模型相关的。我们把实时的、模型无关的特征称为按需特征(on-demand features)。
现在我们已经介绍了数据变换的分类,我们可以深入了解我们的三个 ML 管道,从特征管道开始。
特征管道
特征管道是一个程序,它编排模型无关和按需数据变换的数据流图(dataflow graph)的执行。这些变换包括从数据源提取数据、数据验证和清洗、特征提取、聚合、降维(例如创建向量嵌入)、分箱(binning)、特征交叉(feature crossing),以及对输入数据进行其他特征工程步骤,以创建和/或更新特征数据。
批量或流式特征管道可以应用这些类型的数据变换中的部分或全部,以创建存储在特征存储中的特征,如 图 2-7 所示。该图还显示了另外两种专门的特征管道:向量嵌入管道,它创建存储在(特征存储中的)向量索引(vector index)里的向量嵌入;以及特征数据验证管道,它是一个异步程序,对存储在特征存储中的特征数据运行数据验证规则。
显示三种特征管道类型的示意图:批量或流式特征管道将数据从数据源变换为特征,向量嵌入管道创建嵌入,特征数据验证管道对特征数据执行异步验证。

然而,特征管道不仅仅是一个执行数据变换的程序。它必须能够连接数据源并从中读取数据,需要将特征数据保存到特征存储,而且还有非功能性需求,例如:
- 回填或增量数据
- 同一个特征管道应该能够使用历史数据或新数据(生产数据)创建特征数据,新数据以批量形式到达或作为传入数据流到达。
- 容错
- 特征管道中的失败和重试不应导致数据损坏或重复。
- 可扩展性
- 你必须确保特征管道配置了足够的资源来处理预期的数据量。
- 特征新鲜度
- 客户端使用的预计算特征数据允许的最大年龄是多少?特征新鲜度要求是否意味着你必须将特征管道实现为流处理程序,还是可以是批量程序?
- 治理与安全要求
- 数据可以在哪里处理、谁可以处理数据、处理是否会创建防篡改的审计日志,以及特征是否会为可发现性而组织和打标签?
- 数据质量保证
- 你的特征管道是在数据写入特征存储之前验证数据,还是在数据落入特征存储之后异步验证?
让我们从特征管道的数据源开始——它们从哪里来?想象一下开发一个新的特征管道,并从你从未解析过的数据源获取数据(例如,数据仓库中的现有表)。该表可能已经积累了一段时间的数据,因此你可以对表中的历史数据运行数据变换,以将特征数据回填(backfill)到特征存储中。你可能还想更改特征管道中的数据变换,因此你再次需要(使用新的变换)从源表回填特征数据。你的数据仓库表可能还会以某种节奏(例如每小时或每天)产生新数据。在这种情况下,你的特征管道应该能够从表中提取新数据,计算新的特征数据,并对特征存储中的特征数据进行增量更新(追加、删除或更新)。
你的特征管道创建的特征数据是什么样的?输出的特征数据通常是表格格式(一个或多个 DataFrame 或表),并且通常存储在特征存储中的特征组中。特征组将特征数据存储在表中,客户端在训练和推理(在线应用和批量程序)中都使用这些表。
理想情况下,特征管道应该通过幂等(idempotent)和对特征组的原子更新来容忍失败。幂等性意味着即使运行多次,它们也应该产生相同的结果。原子性(atomicity)意味着更新应该一次性全部应用,因此如果特征管道在完成前失败,带有损坏或缺失数据的部分更新就不应该应用到特征组。幂等性和原子性的好处是,在发生故障时你可以安全地重新运行特征管道。
你可以通过在许多不同的框架和语言之一中实现特征管道来满足可扩展性和特征新鲜度要求。你必须根据特征新鲜度要求、数据输入规模以及团队可用技能来选择最佳技术。不同的数据处理引擎在(1)高效处理、(2)可扩展处理以及(3)开发和运维的易用性方面具有不同的能力。例如,如果批量特征管道每次执行处理的数据少于 1 GB,Pandas 是一个很好的起步框架。对于数十 GB 的工作负载,Polars 是一个不错的选择。对于 TB 级工作负载,Apache Spark 和 SQL 是流行的选择。虽然我们目前介绍的是 DataFrame 处理框架,但 dbt 也是执行用 SQL 定义的特征管道的流行框架。dbt 框架通过允许将变换定义为单独的文件(dbt 称之为模型)作为一种管道形式,为 SQL 增加了一些模块化。这些管道可以链接在一起实现特征管道,最终输出到特征存储中的特征组。
当你的 AI 系统需要新鲜的特征数据时,你可能需要使用流处理来计算特征。对于产生最新鲜特征的流处理特征管道,Feldera 是一个开源的基于 SQL 的引擎,入门门槛低,而更大规模的工作负载可以使用 Apache Flink,它可以扩展到 PB 级工作负载。如果你想要基于 Python 的流式特征管道,那么 Spark Structured Streaming 是一个合理的选择,尽管由于它以批处理方式处理事件(而不是逐事件处理),它会比 Feldera 或 Flink 引入更多的延迟。我们将在第8章中介绍批量特征管道,并在第9章中介绍流式特征管道。
训练管道
训练管道是一个程序,它执行从特征存储读取特征数据、对特征数据应用模型相关变换、用 ML 框架训练模型、验证训练好的模型的性能和没有偏差、将模型发布到模型注册表,最后将模型部署到生产环境进行推理等任务。训练管道按需运行或按计划运行(例如,新模型可以每天或每周重新训练和重新部署一次)。
图 2-8 显示了四类训练管道。第一类是完整的训练管道,它执行所有训练管道任务。它从特征存储中选择、过滤和连接所需的特征数据开始,在将训练好并验证过的模型上传到模型注册表时完成。
说明四类训练管道的示意图:训练、模型验证、模型部署和训练数据集,详细说明特征选择、验证和部署等具体任务。

其他专门的训练管道可以执行这些任务的子集。模型部署管道从模型注册表下载模型,并将其部署用于批量或在线服务。对于在线模型,模型通常部署到模型服务基础设施。它通常是与训练管道分开的管道,因为它是一个运营步骤,可能需要人工批准,并且如果部署出现问题可能需要回滚。模型部署通常涉及A/B 测试(A/B test),其中模型首先作为影子版本(shadow version)部署,如果它表现出足够好的性能和表现,之后被提升为活跃版本。
模型验证也可以在自己的模型验证管道中执行,在模型保存到模型注册表后,对其性能和合规性进行异步评估。当模型验证是计算密集型步骤、不需要 GPU,而模型训练管道使用 GPU 时,这很有用。这样,模型训练可以完成并释放 GPU,模型验证可以在之后用更便宜的 CPU 运行。
对于训练数据集较大且需要时间物化(materialize)的模型,训练管道可以进一步分解为训练数据集管道,它从特征存储中选择、过滤和连接特征数据,对特征数据应用模型相关变换,并将最终训练数据保存为文件。这些文件存储在文件系统或对象存储(如 S3)中。
推理管道
推理管道是一个程序,它读取新的特征数据(预计算的或作为预测请求中的参数),对特征数据应用变换(按需和/或模型相关变换),并用模型输出预测。根据 ML 系统是实时(交互式)ML 系统还是批量 ML 系统,你的推理管道将是批量程序,或是由模型服务基础设施上的预测请求调用的程序。智能体大多是交互式 AI 系统,客户端查询会触发 LLM 调用和外部工具执行的循环,然后才返回响应。
图 2-9 显示了三类推理管道:批量推理管道、在线推理管道和智能体管道。
说明三类推理管道的示意图:批量推理、在线推理和智能体管道,显示每种类型的数据和模型处理步骤流程。

批量推理管道从特征存储中读取预计算特征形式的推理数据,从模型注册表下载模型,并以推理数据作为模型输入输出预测。批量推理管道通常实现为使用 Pandas、Polars 或 Spark 的 Python 程序,尽管一些数据仓库现在支持使用 SQL 的 UDF 进行批量推理。批量推理管道由某个编排器(如 Apache Airflow)按计划运行,使用模型为输入 DataFrame(或 SQL 表)中的所有行进行预测,预测结果通常存储在数据库的表中,消费者从那里使用这些预测。批量推理 ML 系统的一个例子,是我为爱尔兰一处海滩(拉欣奇,Lahinch,我在那里冲浪过很多次)编写的每日冲浪高度预测服务。它从天气和海洋涌浪预报网站抓取数据,并每天发布一个仪表板。批量推理管道往往没有大量参数。也许它们会以推理数据的 start_time 和 end_time 或 last_processed_timestamp 为参数。或者推理数据可能是一组实体(例如用户),在这种情况下,我们将实体 ID 作为参数传递。
在线推理管道获取请求参数,在需要时从特征存储读取任何预计算特征,对预计算特征和请求参数执行任何数据变换以创建特征向量,用特征向量调用模型,记录预测和特征(用于监控和调试),最后将预测返回给客户端。在线推理管道是一种网络托管的服务,响应预测请求进行预测。它通常是一个 Python 程序,有一个称为部署 API(deployment API)的 API,在第11章中描述。部署 API 包括正在为其进行预测的实体的 ID,以及计算模型实时特征所需的任何参数。ID 用于检索实体的预计算特征。在线推理管道的 Python 程序通常与模型一起部署在模型服务基础设施上,例如 KServe 或 FastAPI 服务器。详见第11章。
最后,智能体管道(agentic pipeline)与在线推理管道类似,它是一个网络托管的 Python 程序,具有用于客户端查询和响应的部署 API。智能体本身通常用智能体框架编写,例如 LlamaIndex、LangGraph、LangChain 或 CrewAI。智能体程序有一个 LLM,并且可以访问一组工具以及每个工具的模式(schema)。工具(tool)是智能体可以执行的动作(例如进行外部 API 调用或查询检索增强生成(retrieval-augmented generation,RAG)数据源)。可用工具集要么是静态定义的,要么由智能体在运行时发现。智能体收到客户端的查询后,它会循环执行以下操作,直到向客户端返回响应。首先,它将查询和可用工具列表发送给 LLM,并询问 LLM 应该使用哪个工具。LLM 返回一个或多个要使用的工具及其参数,或返回给客户端的最终响应。如果 LLM 回复要使用的工具,智能体执行该工具,并将结果连同之前的工具使用历史一起发送给 LLM。智能体持续循环进行工具使用/响应步骤,直到 LLM 指示应向客户端发送最终响应。
用 ML 管道构建的泰坦尼克号生存 ML 系统
我们现在介绍我们的第一个示例 ML 系统,它是用我们的三个 ML 管道构建的,使用最著名的 ML 问题之一——预测乘客在泰坦尼克号(Titanic)沉没中幸存的概率。泰坦尼克号乘客生存数据是一个静态数据集。ML 模型在静态数据集上训练和评估。这使它成为学习 ML 的好入门数据集,因为你跳过了创建训练数据的步骤。但我们想超越仅仅用静态数据转储训练模型的想法。
在 图 2-10 中,我们在看板上看到 ML 系统的轮廓,包括其数据源、最终输出(一个仪表板)以及用于实现 ML 系统的技术。
泰坦尼克号乘客生存 ML 系统的 MVPS 看板示意图,显示数据源、ML 管道和作为生存仪表板的最终输出。

我们将使用泰坦尼克号生存数据集作为历史数据,如 图 2-11 所示。
泰坦尼克号生存数据集表,显示乘客信息和结果,突出显示用于机器学习的特征和用于监督学习任务的标签。

然后我们将编写一个合成数据创建函数,为泰坦尼克号创建新乘客。合成乘客的特征值从与原始数据集相同的分布中抽取,因此我们不会遇到特征漂移(feature drift)问题,也不需要重新训练模型。这是一个过于简化的例子,但对于开始使用动态数据仍然很有用。
我们将使用 Pandas 编写的 Python 特征管道,将历史特征数据和新特征数据写入 Hopsworks 特征存储中的单个特征组。然后我们将安排特征管道每天运行一次,为泰坦尼克号每天创建一个新乘客:
import pandas as pd
import hopsworks
BACKFILL=True
def get_new_synthetic_passenger():
# see github repo for details
if BACKFILL==True:
df = pd.read_csv("titanic.csv")
# Remove columns that are not predictive of passenger survival
else:
df = get_new_synthetic_passenger()
fs = hopsworks.login().get_feature_store()
fg = fs.get_or_create_feature_group(name="titanic", version=1, \
primary_keys=['id'], description="Titanic passengers")
fg.insert(df)
我们的训练管道首先选择我们想在模型中使用的特征,并创建一个特征视图(feature view)来表示模型的输入特征和输出标签/目标。我们使用特征视图从泰坦尼克号乘客生存数据中读取训练数据,这些数据被随机分为 20% 的测试集和 80% 的训练集。然后我们将用 XGBoost(一个 Python 中的梯度提升决策树库)训练模型。最后,我们将训练好的模型保存到 Hopsworks 的模型注册表:
import xgboost
fg = fs.get_feature_group(name="titanic", version=1)
fv = fs.get_or_create_feature_view(name="titanic", version=1, \
labels=['survived'], \
query=fg.select_features()
)
X_train, X_test, y_train, y_test = fv.train_test_split(test_size=0.2)
model = xgboost.XGBClassifier()
model.fit(X_train, y_train)
model.save_model("model_dir/model.json")
mr = hopsworks.login().get_model_registry()
mr_model = mr.python.create_model(
name="titanic",
feature_view=fv,
)
mr_model.save("model_dir")
我们将编写一个计划每天运行一次的批量推理管道。它将从特征存储中读取任何新的合成乘客,从模型注册表下载我们训练好的模型,使用模型预测合成乘客是否幸存,并将预测与特征视图一起记录到特征存储:
retrieved_model = mr.get_model(name="titanic", version=1)
saved_model_dir = retrieved_model.download()
model = xgboost.XGBClassifier()
model.load_model(saved_model_dir + "/model.json")
row_data = # get row of features for new passenger
prediction = model.predict(row_data)
这个 ML 系统解决的是所谓的反事实(假设情景)预测问题(counterfactual(what-if)prediction problem)。如果泰坦尼克号上有一位男性乘客,49 岁,乘坐三等舱——他幸存的概率是多少?这个「泰坦尼克号乘客生存作为 ML 系统」示例的完整源代码可以在 GitHub 上的本书源代码仓库中找到。它还包括一个用 Python 编写的交互式 UI(使用 Gradio),用于向模型询问有关乘客生存概率的假设情景问题。
要开始使用这个示例,你需要安装 Hopsworks Python 客户端库。在 Linux 和 Apple 系统上,这涉及调用:
pip install hopsworks[python]
在 Windows 上,你需要先 pip install Twofish 库,然后才能安装 Hopsworks 库。你还需要在 Hopsworks Serverless 上创建一个账户,并且需要获取一个 Hopsworks API 密钥(用户 → 账户 → API),并将其保存到课程仓库根目录的 .env 文件中,以便你可以安全地从 Hopsworks 读取和写入。你可以运行第一个笔记本,让它提示你创建 Hopsworks API 密钥,或者你可以按照文档操作。Hopsworks 提供永久免费的 serverless 套餐,包含 35 GB 的免费存储空间,这足以完成本书中的项目。
小结
在构建 AI 系统时,我们从 ML 管道以及特征管道、训练管道和推理管道中执行的数据变换开始。我们介绍了 ML 管道数据变换的分类法,其基础是可复用特征(由特征管道中的模型无关变换创建)、模型专属特征(由训练/推理管道中的模型相关变换创建)和实时特征(由在线推理管道中的按需变换创建,也可以应用于历史数据以在特征管道中创建特征)。在本章的最后,我们介绍了第一个 ML 系统——泰坦尼克号乘客生存预测问题的动态数据版本。我们展示了如何为泰坦尼克号乘客生存构建批量 ML 系统和交互式 ML 系统。在下一章中,我们将更进一步,你将为你所在的街区或地区构建一个 AI 系统。你将为你居住的街区构建一个空气质量预测服务,我们将使用与泰坦尼克号示例相同的框架——Python、Pandas、XGBoost 和 Gradio。
1 Edsger W. Dijkstra,《Go To 语句被认为有害》(“Go To Statement Considered Harmful”),Communications of the ACM(《ACM 通讯》)第 11 卷,第 3 期(1968 年 3 月):147-48。