你身边的空气质量预测服务

我们要构建的第一个 ML 项目,是为一个你在意的社区提供空气质量预测服务。我们将遵循第 2 章的 MVPS 流程——divide et impera(分而治之)。你的成果将是一项面向公众、旨在长久存续的服务,所以请在上面多花些时间和心思,你的社区会因此感谢你。我对这个项目有个人感情,因为我有两个患囊性纤维化(cystic fibrosis)的儿子——这是一种主要影响肺部的遗传性疾病。他们出生在同一天,相差两岁,也在同一天被确诊。总之,我认为我可以代表整个囊性纤维化社区说:对我们和许多人来说,这都会是一项很棒的服务!1

我们的 AI 系统要解决的预测问题是:预测你家、公司或任何地方附近的一个公共空气质量传感器所测得的空气质量。全球有一个由物联网(Internet of Things,IoT)爱好者组成的社区,他们把传感器装在自己的花园和阳台上,并把空气质量测量数据发布到互联网上。我居住的斯德哥尔摩有 30 多个公共传感器,而我的家乡都柏林有 40 多个。有一个世界空气质量指数网站,你可以在它的地图上找到一个传感器,用来构建你的 AI 系统。选一个同时满足以下两个条件的传感器:(1) 有历史数据——我们要在历史数据上训练 ML 模型,所以如果你有几年数据,那就太好了;(2) 能产出可靠的测量值(有些传感器会在一段时间内关机或出现故障)。可靠的传感器能让你的 AI 系统持续收集测量数据,随着更多数据的积累,你就可以重新训练并改进模型。尽管你将为社区提供一项免费的公共服务,但你一分钱也不用花——我们将把系统运行在免费的 Serverless 服务(GitHub 和 Hopsworks)上。

空气质量预测是一个相当直接的 ML 问题。我们把这个预测问题建模为回归(regression)问题——预测 PM 2.5 的值。PM 2.5 是一种细颗粒物指标,衡量直径小于等于 2.5 微米的颗粒物,其浓度过高会增加健康问题的风险,比如低出生体重、心脏病和肺病。PM 2.5 浓度过高还会降低能见度,使空气看起来灰蒙蒙的。我们将用什么特征来预测 PM 2.5 的水平呢?PM 2.5 与风速/风向、温度和降水相关,所以我们将使用天气预报数据来预测以 PM 2.5 衡量的空气质量。这很有道理,因为当风从特定方向吹来时,空气质量通常会更好——如果你住在繁忙的道路旁边,风向就至关重要。在较冷的天气里空气质量往往更差,因为冷空气比暖空气更稠密、移动更慢;而在城市里,通勤时开车的人可能比骑自行车的人更多。即使是印度那些冬季并不寒冷的地区,冬季月份的空气质量也更差。

但是,等等。你可能读到过,空气质量预测已经是一个被解决的问题。2024 年,微软 AI 构建了 Aurora,一个能预测全球空气污染的深度学习模型。与欧盟的哥白尼计划在高性能计算基础设施上计算的空气物理模型相比,微软对 AI 的运用被誉为一大进步。然而,截至 2024 年年中,如果你在某个城市(比如斯德哥尔摩)检验 Aurora 的表现,你会发现它的预测与你可以在 aqicn.org 上找到的实际空气质量传感器读数相比并不准确。因此,你的挑战是构建一个 AI 系统,在你选择的空气质量传感器所在位置,用其成本的一小部分,就产生比 Aurora 更好的空气质量预测。在这个项目中,更高质量的数据加上决策树 ML 模型将胜过深度学习。

最后,每个项目都会受益于一个"哇"的因素。我们将给项目撒上一些生成式 AI(GenAI)的魔法,通过一个由 LLM 驱动的语音 UI,让你的空气质量预测服务变得"友好"。

AI 系统总览

在我的 KTH 课程中,学生以项目形式构建一个独特的 AI 系统,用动态数据源解决一个预测问题。但在开始项目之前,他们必须先获得批准,我发现最简单的批准方式就是预测服务卡片(见表 3-1)。这张卡片是图 2-2中看板(kanban board)的简化版,省略了实现细节。

动态数据源预测问题UI 或 API监控
空气质量传感器数据:https://aqicn.org 天气预报:https://open-meteo.com在某个现有空气质量传感器的位置,对未来七天的 PM 2.5 水平进行每日预测包含图表的网页,以及由 LLM 驱动的 Python UI后报(hindcast)图展示我们模型的预测表现

AI 系统卡片简明扼要地总结了系统的关键属性,包括数据源和它要解决的预测问题。例如,就空气质量而言,存在许多可能的空气质量预测问题,比如预测 PM 10 水平(更大的颗粒物,包括来自道路和建筑工地的灰尘)和 NO 2(二氧化氮)水平(主要来自内燃机车辆的污染)。预测服务卡片还包含数据源,这使它可以用作可行性测试,检验数据是否存在且对你的预测问题可访问。你还应该定义 AI 系统产生的预测将如何被消费——通过 UI 或 API。UI 是向利益相关者传达模型价值的非常强大的工具,而现在用 Python 构建可用的 UI 已经很直接了。在我们的 AI 系统中,我们将使用 LLM 来提高服务的可访问性——你应该能用自然语言向空气质量预测服务提问。最后,你应该概述如何监控运行中的 AI 系统的表现,以确保它按预期运行。

我们将使用开源和免费的 Serverless 服务来构建我们的 AI 系统——GitHub Actions/Pages 和 Hopsworks。我们将用 Python 编写以下 Jupyter notebooks:

  • 用于存储我们的数据并用历史数据回填(backfill)的特征组
  • 一条每日特征管道(feature pipeline),用于检索新数据并将其存储在特征存储(feature store)中
  • 一条训练管道(training pipeline),用于训练 XGBoost 回归模型并将其保存到模型注册表(model registry)
  • 一条批推理管道(batch inference pipeline),用于下载模型、对新特征数据做预测,并从特征存储读取数据,以生成空气质量预报/后报图

我们还将使用一些 Python 库和其他技术来构建该系统,包括:

  • 用于从我们的空气质量和天气数据源读取数据的 REST API
  • 用于处理数据的 Pandas
  • 用于存储特征数据和模型的 Hopsworks
  • 用于我们 ML 模型(梯度提升决策树)的 XGBoost
  • 用于调度我们的 notebooks 每日运行的 GitHub Actions
  • 作为包含预报/后报图的仪表板网页的 GitHub Pages

我们还将编写一个 Streamlit Python 应用,提供语音和文本驱动的 UI,背后是开源的 Whisper transformer 模型(把语音转成文本)和一个 LLM(把文本转成对我们 AI 系统的函数调用)。

对我们第一个项目来说,技术有点多,但别被吓到。就像很多伟大的音乐可以用三个和弦谱写出来一样,很多伟大的 AI 系统可以由一条特征管道、一条训练管道和一条推理管道构成。

空气质量数据

世界各地有成千上万的爱好者安装了空气质量传感器,并把他们的测量数据公开免费地分享出来。你可以用 aqicn.org 地图找到许多同时具有历史数据和实时数据的空气质量传感器。这个网站是来自许多来源的传感器数据的聚合器,但作为一项社区服务,它不对数据质量作任何保证。

我选择了一个斯德哥尔摩的传感器,它同时提供实时和历史数据(见图 3-1)。我选择它是因为它离 Hopsworks 的办公室非常近。

热力图展示了一个斯德哥尔摩传感器多年来的空气质量数据,并醒目地高亮显示了"Download this data (CSV format)"(以 CSV 格式下载此数据)按钮。

原书插图

你应该选择一个离你近的,或者对你来说有特殊意义的传感器。向下滚动页面,你会找到一个下载该传感器历史数据的按钮。如果你在传感器网页上找不到历史测量数据的下载链接,你也许能在世界空气质量历史数据库里找到。如果还是找不到下载链接,就换一个传感器。遗憾的是,截至 2025 年年中,还没有可用的 API 调用来下载历史数据,所以这一步你必须手动完成。你还需要在 AQICN 网站上创建一个 API 密钥,这样你的特征管道才能读取最新的空气质量值。

下载 CSV 文件。我把我的文件重命名为 air-quality-data.csv。对于你的传感器,如果你下载的 CSV 文件名里有空格或不寻常的字符,你应该重命名它。你应该用文本编辑器打开 CSV 文件,检查列名是否符合预期。我们的回填 Python 程序会把 CSV 文件读入 Pandas DataFrame,并期望 CSV 文件有一个标题行,且其中两列名为 pm25date。如果有更多列也没关系,程序会忽略它们。但是,有些文件没有 pm25 列——取而代之的是 PM 2.5 的 min/max/median/stdev 每日测量值。最简单的修复方法是直接把 CSV 文件标题中的 median 列改名为 pm25。你还必须有 date 列。

现在你可以通过把本书的 GitHub 仓库 fork 到你的 GitHub 账户来创建项目的 GitHub 仓库。你应该把你的 CSV 文件移动到 fork 后仓库的 data 目录,并替换现有的 data/air-quality-data.csv 文件。你还应该从 .env.example 模板创建一个 .env 文件。你需要在 .env 文件中用你的 API 密钥值,以及你所选传感器的 URL、国家、城市和街道,更新以下值:

HOPSWORKS_API_KEY=< get your key from Hopsworks > AQICN_API_KEY=< get your key from aqicn.org > AQICN_URL=https://api.waqi.info/feed/@10009 AQICN_COUNTRY=sweden AQICN_CITY=stockholm AQICN_STREET=hornsgatan-108

.env 文件不应该提交到 GitHub(它已在 .gitignore 文件中)。把你的 CSV 文件提交并推送到 GitHub。CSV 文件相当小(我的只有 58 KB),所以把它们存在 GitHub 上没有问题。GB 级或更大的文件不适合存储在像 GitHub 这样的源代码仓库中。2 当你在 Python 中工作时,我们强烈建议你为这本书创建一个虚拟环境(virtual environment),使用诸如 CondaPoetryvirtualenvpipenv 之类的 Python 依赖管理框架。我们项目引入的依赖可以安装在你的虚拟环境中。关于为这个项目设置虚拟环境和安装 Python 依赖的细节,请参阅本书的源代码仓库。在第 2 章中,我们已经讨论过如何创建你的 Hopsworks 账户并下载 API 密钥。

探索性数据集分析

在我们开始构建之前,应该花点时间了解我们将要处理的数据。一般来说,在把任何数据源用于解决预测问题之前,你应该了解它的六个属性或维度:

  • 有效性(Validity)
  • 准确性(Accuracy)
  • 一致性(Consistency)
  • 唯一性(Uniqueness)
  • 更新频率(Update frequency)
  • 完整性(Completeness)

现在让我们用这个透镜来审视我们的空气质量和天气数据源。

Note

我们建议本书使用 Jupyter Notebooks 而不是 Google Colaboratory(Colab)。每个 notebook 的第一个单元都添加了对 Colab 的支持,但你必须更新它以指向你 fork 的仓库。Colab 目前对 GitHub 的支持不太好,所以每个 notebook 都必须先克隆仓库并安装所有依赖才能运行。而且不支持把你对 notebook 所做的任何更改保存回 GitHub。不过,如果你需要免费的 GPU,Colab 仍然有用。

空气质量数据

我们的空气质量数据源在这六个数据集质量属性上表现如何?

我们从数据有效性(data validity)开始,它衡量数据在多大程度上准确地反映了它要测量的对象。我们专注于测量 PM 2.5 而不是 PM 10 或 NO 2,因为根据联合国的说法,按照目前的认识,“PM 2.5……构成最大的健康威胁”。接下来是数据准确性(data accuracy),它指的是测量值与真实值的接近程度。aqicn.org 网站告诉我,我的斯德哥尔摩传感器的数据来自"SLB·analys——斯德哥尔摩市空气质量管理与运营机构"和"欧洲环境署"。因此,我愿意信任这些数据的准确性。

回到 stockholm-hornsgatan-108 数据集,我们声称数据是唯一的(unique)。经过网络搜索,我不知道那条街上还有其他任何公共空气质量传感器。看看图 3-1中的数据,我可以看到数据基本完整、相当一致(consistent)(表示空气质量的颜色遵循可预测的模式),并且是及时的(每小时到达一次)。

一般来说,你还应该在 notebook 中检查数据的完整性(completeness)。在下面的代码片段中,我们把 CSV 文件读成 Pandas DataFrame,然后只保留空气质量数据集中需要的那些列(date 和我们的目标 pm25):

# you may need to rename columns in your CSV file to 'pm25' and 'date'
df = pd.read_csv("../../data/stockholm-hornsgatan-108.csv",
parse_dates=['date'], skipinitialspace=True)
df_aq = df[["date", "pm25"]]

我们还使用 Pydantic settings 对象从 .env 中读取传感器的 countrycitystreeturl,并把它们作为列添加到 df_aq。我们将使用 city 列把空气质量数据与同一 date 的天气特征连接起来。我们用 city 值来检索下载天气数据所需的经度纬度countrycitystreet 列是辅助列(helper columns),在我们创建带空气质量预报的仪表板时会用到。我们还将 countrycitystreeturlHOPSWORKS_API_KEYAQICN_API_KEY 作为 secret 存储在 Hopsworks 中,这样后面的 notebooks(每日特征管道、训练管道和推理管道)就不需要从 .env 文件中读取这些值了。

评估数据集完整性的第二部分是检查缺失数据。你可以在 DataFrame 上调用 isna() 函数来列出任何缺失值。但是,那可能会输出大量行,所以我们将对 isna() 的结果应用 sum(),汇总 df 中每列缺失值的数量:

df.isna().sum()

然后你可以通过调用以下代码删除任何包含缺失列的的行:

df.dropna(inplace=True)

在这个阶段删除缺失的观测值是合理的,因为收集日期或目标缺失的数据没有意义。

通常,到这个时候我们会更深入地识别数据源和模型的候选特征。我们会尝试识别对目标(PM 2.5)具有预测能力的特征。如果没有足够的样本让深度学习模型表现出色,我们可能会尝试工程化一些能捕捉预测问题领域知识的特征。但是,在这种情况下我们将跳过这些步骤,把它建模为一个更简单的预测问题。我们将为模型使用天气特征,因为它们对 PM 2.5 水平有很好的预测能力。我们要训练的模型还有改进空间,但现在,我们的目标是为空气质量预测问题构建一个 MVPS。

天气数据

我们将使用 Open-Meteo 下载与你所选空气质量传感器相同位置的历史天气数据和天气预报数据。Open-Meteo 的天气数据在我们数据质量的六个维度上得分都很高。Open-Meteo 提供两个不同的免费 API:一个用于下载历史天气数据,一个用于天气预报。你不需要 API 密钥。如果你不确定哪个城市最适合你的天气数据,你可以在 Open-Meteo 历史天气 API 页面搜索可用的天气位置。与高度本地化的空气质量数据(相邻两条街的空气质量状况可能截然不同)相比,城市甚至区域级别的天气数据对你的模型来说可能就足够了。

我们将把自己限制在那些天气站普遍可用、且对空气质量预测能力最强的天气条件上:降水、风速、风向和温度。Open-Meteo API 期望以经度和纬度作为你天气位置的参数。我们使用 geopy 库来解析你指定城市名的经度和纬度(如果 geopy 服务器屏蔽了你的 IP,你可能需要手动输入经度和纬度)。

在下面使用历史 API 的代码片段中,我们需要把位置和时间范围作为 longitudelatitudestart_dateend_date 参数提供:

url = "https://archive-api.open-meteo.com/v1/archive" params = { "latitude": latitude , "longitude": longitude , "start_date": start_date , "end_date": end_date , "daily": ["temperature_2m_mean", "precipitation_sum", "wind_speed_10m_max", "wind_direction_10m_dominant"] } responses = openmeteo.weather_api(url, params=params)

天气预报数据将通过类似的 REST 调用检索:

url = "https://api.open-meteo.com/v1/ecmwf" params = { "latitude": latitude , "longitude": longitude , "daily": ["temperature_2m", "precipitation", "wind_speed_10m", "wind_direction_10m"] } responses = openmeteo.weather_api(url, params=params)

但是,你应该注意,我们的预报 API 调用接收的是逐小时的预报,而我们的历史 API 调用检索的是一天内的聚合数据(即平均温度、降水总量和最大风速)。这并不理想,但对我们的目的来说已经足够好了(我们说过模型还可以改进!)。

weather-util.py 中定义了两个实用函数,get_historical_weather()get_weather_forecast(),它们以 Pandas DataFrame 的形式返回天气数据:

historical_weather_df = util.get_historical_weather("Stockholm", "2019-01-01", 
    "2024-03-01")
weather_forecast_df = util.get_weather_forecast("Stockholm")
Warning

注意,这些函数会发起网络调用,所以如果程序没有互联网连接,代码可能会失败。我们将用来检索实时空气质量数据的函数也是如此。

创建并回填特征组

我们将把特征化(featurized)的 DataFrames 存储在 Hopsworks 特征存储的特征组(feature group)中。我们会有两个特征组:一个用于空气质量数据(包含 PM 2.5 值的观测、位置以及这些观测的时间戳),另一个用于存储历史天气观测和天气预报数据。特征组存储随时间累积的增量特征数据:

air_quality_fg = fs.get_or_create_feature_group(
    name='air_quality',
    description='Air Quality observations daily',
    version=1,
    primary_key=['country', 'city', 'street'], 
    expectation_suite = aq_expectation_suite,
    event_time="date",
)    
air_quality_fg.insert(df_aq)

我们调用 get_or_create_feature_group() 而不是 create_feature_group(),因为我们希望 notebook 是幂等的(如果特征组已经存在,create_feature_group() 会失败):

weather_fg = fs.get_or_create_feature_group(
    name='weather',
    description='Historical daily weather observations and weather forecasts',
    version=1,
    primary_key=['city'], 
    event_time="date",
    expectation_suite = weather_expectation_suite
) 
weather_fg.insert(df_weather)

注意,两个特征组都定义了一个 expectation_suite 参数。这是一组数据验证规则,我们以声明式的方式一次性附加到特征组上,但每次向特征组写入 DataFrame 时都会执行。

我们可以定义数据质量测试来验证从空气质量和天气数据源检索到的数据。这些测试将从故障一开始发生就能帮助识别传感器中的故障。Great Expectations 是一个流行的开源库,用于以声明式方式指定数据验证规则。在下面的代码片段中,我们在 Great Expectations 中定义了一个期望(expectation),检查我们 DataFrame dfpm25 列的所有值,确保抓取到的值既不是负数也不大于 500(对我所在位置的 PM 2.5 预期值来说,这是一个合理的上限):

import great_expectations as ge
aq_expectation_suite = ge.core.ExpectationSuite(
    expectation_suite_name="aq_expectation_suite"
)

aq_expectation_suite.add_expectation(
    ge.core.ExpectationConfiguration(
        expectation_type="expect_column_min_to_be_between"
        kwargs={
            "column":"pm25",
            "min_value":0.0,
            "max_value":500.0,
            "strict_min":True
        }
    )
)

在 Hopsworks 中,如果数据验证规则失败,你可以轻松地添加通知(通过 Slack 或电子邮件),并把策略设置为要么摄取数据并警告,要么使摄取失败。在本书的源代码仓库中,还为天气数据的 temperature_2mprecipitation 列定义了类似的期望。

特征管道

我们刚才展示了创建特征组并用历史数据回填它们的程序。但我们也需要每天处理新数据。我们可以扩展之前的程序,并参数化它以在回填模式或正常模式下运行。但是,相反,我们将把每日特征管道写成单独的程序——这样就把创建特征组和回填的职责与特征组的每日更新分开了。回填和每日特征管道使用的公共函数定义在 mlfs/airquality 包的模块中。

每日特征管道将被调度为每天运行一次,执行以下任务:

  • 读取今天的 PM 2.5 测量值
  • 读取今天的天气数据测量值
  • 读取未来七天的天气预报数据
  • 把所有这些数据插入空气质量和天气特征组

这个例子不需要特征工程。我们将把所有数据读取为数值特征数据,并且在写入特征组之前不会对数据进行编码。下载传感器读数和天气预报的代码位于 functions/util.py 模块中:

url = f"{aqicn_url}/?token={ AQI_API_KEY `}" data = trigger_request(url) aq_today_df = pd.DataFrame() aq_today_df[‘pm25’] = [data[‘data’][‘iaqi’].get(‘pm25’, {}).get(‘v’, None)] aq_today_df[‘city’] = city .. aq_today_df[‘date’] = datetime.date.today() air_quality_fg.insert(df_air_quality)

url = “https://api.open-meteo.com/v1/ecmwf" params = { “latitude”: latitude, “longitude”: longitude, “hourly”: [“temperature_2m”, “precipitation”, “wind_speed_10m”, “wind_direction_10m”] } responses = openmeteo.weather_api(url, params=params) hourly_df = # populate with responses data daily_df = hourly_df.between_time(‘11:59’, ‘12:01’) weather_fg.insert(daily_df)`

我们对 aqicn 和 Open-Meteo 的 API 调用分别返回空气质量和天气预报数据,我们把返回的数据放入 Pandas DataFrames,然后插入到各自的特征组中。当你把 DataFrame 插入特征组时,它的数据验证规则将被执行。

你可以在 Hopsworks UI 中看到历史特征管道执行的结果。登录 Hopsworks,导航到"Feature group”→“Recent activity”(特征组 → 近期活动)查看摄取运行的结果。你可以在"Feature group"→“Data preview”(特征组 → 数据预览)中检查特征组的内容。在"Feature group"→“Feature statistics”(特征组 → 特征统计)中查看对插入数据计算出的描述性统计,并在"Feature group"→“Expectations”(特征组 → 期望)中查看数据验证结果。

训练管道

我们决定把 PM 2.5 建模为回归问题,而且我们知道我们只会有几百行,可能一千行左右的数据。这绝对属于小数据的范畴,所以我们不会使用深度学习。相反,我们将使用小数据的首选 ML 框架——XGBoost,一个开源的梯度提升决策树框架。XGBoost 开箱即用表现良好,我们在这里不做任何超参数调优——我们把它留给你作为练习,从模型中榨取更多性能。

我们将从选择要在模型中使用的特征开始。为此,我们将使用 Hopsworks 中的特征视图(feature view)。特征视图定义了一个模型的模式(schema)——它的输入特征和输出目标(或标签)。Hopsworks 提供了一个类似 Pandas 的 API,用于从不同特征组中选择特征,然后使用 query 对象把选中的特征连接在一起。特征组上的 select()select_all() 方法返回一个 query 对象,该对象提供了一个 join() 方法(更多细节见第 5 章)。当你创建特征视图时,还要指定选中的特征中哪些是标签列。从特征组中选择特征、使用共同的 'city' 列把它们连接起来、并创建特征视图的代码如下:

selected_features = \
air_quality_fg.select(['pm25']).join(weather_fg.select_all(on=['city']))

feature_view = fs.create_feature_view(
    name='air_quality_fv',
    version=version,
    labels=['pm25'],
    query=selected_features
)

有了特征视图对象,你现在就可以创建训练数据了:

X_train, X_test, y_train, y_test = feature_view.train_test_split(test_size=0.2)

在这里,我们把训练数据读取为 Pandas DataFrames,它们被随机分割(80/20)成训练集特征(X_train)、训练集标签(y_train)、测试集特征(X_test)和测试集标签(y_test)。在一次调用中,train_test_split 读取数据、连接空气质量和天气数据,然后用 Scikit-Learn 把数据随机分割成训练集和测试集的特征与标签。我选择随机分割而不是时间序列分割的原因是,我们选择的特征不依赖时间。一个有用的练习是,通过添加与空气质量相关的特征(历史空气质量、季节性因素等)并改用时间序列分割,来改进这个空气质量模型。

我们现在可以用 XGBoostRegressor 训练我们的模型了。我们只需使用 XGBoostRegressor 的默认超参数,把模型拟合到训练集的特征和标签上:

clf = XGBRegressor()
clf.fit(X_train, y_train)

训练应该只需要几毫秒。然后,你可以用测试集的特征来评估训练好的模型 clf,生成预测 y_pred

y_pred = clf.predict(X_test)
mse = mean_squared_error(y_test, y_pred, squared=False)
r2 = r2_score(y_test, y_pred)
plot_importance(clf, max_num_features=4)

因为我们把 PM 2.5 预测建模为回归问题,所以我们使用均方误差(mean squared error,MSE)和 R 平方(R-squared)误差作为评估模型性能的指标。MSE 的替代方案是平均绝对误差(mean absolute error,MAE),但与 MAE 相比,如果模型的预测与结果相差甚远,MSE 对模型的惩罚更大。使用 scikit-learn 库,当你的结果(y_test)和预测(y_pred)就绪时,只需一次方法调用就能计算许多不同的模型性能指标。我们还计算特征重要性,稍后把它保存为 PNG 文件。

现在,我们需要把这个训练管道的输出——训练好的模型 clf——保存到模型注册表。我们将使用 Hopsworks 模型注册表。这个过程首先把模型保存到本地目录,然后把模型注册到模型注册表,包括它的名称(air_quality_xgboost_model)和描述、它的评估指标,以及用于为模型创建训练数据的特征视图:

model_dir = "air_quality_model"
os.makedirs(model_dir + "/images")
clf.save_model(model_dir + "/model.json")
plt.savefig(model_dir + "/images/feature_importance.png")

mr = project.get_model_registry()
mr.python.create_model(
    name="air_quality_xgboost_model", 
    description="Air Quality (PM2.5) predictor.",
    metrics={ "MSE": mse, "r2": r2 },
    feature_view = feature_view
)
mr.save(model_dir)

模型注册表客户端使用特征视图对象提取模型的模式和血缘(lineage)。包含模型的本地目录中的任何其他文件也会被上传,images 子目录(feature_importance.png)中的任何 PNG/JPEG 文件将显示在"Model evaluation images"(模型评估图)部分(见图 3-2)。

仪表板展示了一个 XGBoost 模型的注册详情,包括 R 平方和 MSE 等指标,以及说明特征重要性和 PM2.5 预测的评估图。

原书插图

注意,每次我们注册模型时,都会得到模型的一个新版本。与特征组和特征视图不同,我们在创建模型时不需要提供版本——新注册的模型会自动分配一个自增的版本号。有了模型注册表中的训练好的模型,我们现在可以编写批推理管道来生成空气质量仪表板了。

批推理管道

批推理管道是一个 Python 程序,它从模型注册表下载训练好的模型,获取天气预报特征数据,并使用模型和天气预报数据预测未来七天的空气质量。我们将做七次不同的预测,七天各一次。我们将使用 Plotly 创建空气质量预报图,把图保存为 PNG 文件,并把 PNG 文件推送到一个包含 GitHub Pages 公共网站的 GitHub 仓库。GitHub Pages 有一个免费层,允许你构建网页、仪表板和个人博客,而且你会为你的网站获得一个专属域名。

首先,我们需要从模型注册表下载模型,并用 XGBRegressor 对象加载它:

model_ref = mr.get_model(
    name="air_quality_xgboost_model",
    version=1,
)

saved_model_dir = model_ref.download()
retrieved_xgboost_model = XGBRegressor()
retrieved_xgboost_model.load_model(saved_model_dir + "/model.json")

然后,我们使用天气特征组读取一批推理数据(我们未来七天的天气预报数据):

batch_df = weather_fg.filter(weather_fg.date >= today).read()

batch_df DataFrame 现在包含未来七天的天气预报特征。有了这些特征,我们现在可以用模型做预测:

features = batch_df[['temperature_2m_mean', 'precipitation_sum', \
    'wind_speed_10m_max', 'wind_direction_10m_dominant']]
batch_df['predicted_pm25'] = model.predict(features)
batch_df['days_before_forecast_day'] = range(1, len(batch_df)+1)

我们把预测存储在 batch_dfpm25_predicted 列中,同时存储预报开始前的天数。共有七个预报,每天一个。第一个提前七天,最后一个提前一天。这个 days_before_forecast_day 列将帮助我们评估模型的性能,取决于它提前多少天进行预报。我们将把 batch_df 保存到特征存储中,用于监控特征/预测,因为 batch_df 包含预测、特征值和辅助列:

monitoring_fg = fs.get_or_create_feature_group(
    name='monitoring_aq',
    description='Monitor Air Quality predictions',
    version=1,
    primary_key=['city', 'street']
)  
monitoring_fg.insert(batch_df)

我们还必须绘制我们的空气质量预测仪表板。我们将使用 plotly 库:

import plotly.express as px
fig = px.line(batch_df , x = "date", y = "pm25_predicted", title = "..")
....
fig.write_image(file="forecast.png", format="png", width=1920, height=1280)

我们将使用一个 GitHub Action 把 forecast.png 文件发布到网页上,如下一节所述(见图 3-3)。

Plotly 图表按时间比较预测与实际的 PM2.5 值,并标注空气质量类别。

原书插图

最后,我们创建一些后报(hindcast)PNG 文件,把我们的模型预测(来自监控特征组数据)与结果(来自空气质量特征组数据)进行比较。详情请参阅本书源代码仓库中的批推理管道 notebook。

运行管道

首先,你应该在笔记本电脑上运行 Jupyter notebooks,确保它们按预期工作。从第一个单元运行到最后一个单元。运行每个 notebook 后,你应该切换到 Hopsworks UI 查看所做的更改——比如创建特征组、写入特征组、创建特征视图、把训练好的模型保存到模型注册表。

首先,运行特征回填 notebook(1_air_quality_feature_backfill.ipynb)。这将创建 air_qualityweather 特征组。然后你应该运行特征管道(2_air_quality_feature_pipeline.ipynb),并检查特征组,看是否有新行按预期添加。然后,你可以通过运行模型训练管道(3_air_quality_training_pipeline.ipynb)来训练模型;验证特征视图(air_quality_fv)已创建,并且训练好的模型在模型注册表中。最后,测试你的批推理管道(4_air_quality_batch_inference.ipynb)是否按预期工作——它应该创建了一个 aq_predictions 特征组。如果你发现 bug,请发布 GitHub issue。如果你能改进代码,请提交拉取请求(pull request,PR)。如果你需要帮助,请在本书 GitHub 仓库中链接的 Hopsworks Slack 上提问。

把管道调度为 GitHub Action

我们将使用 GitHub Actions 来调度特征管道和批推理管道,并使用 GitHub Pages 构建我们的仪表板。截至 2024 年,GitHub 的免费层每月给你 2,000 分钟的免费计算时间。这远远足够运行我们的特征管道和批推理管道了。你可以在笔记本电脑的 Jupyter notebook 上运行训练管道——我们现在不会按计划运行它。对于我们的 UI,我们将使用 GitHub Pages(它托管你 GitHub 仓库的网页);在 GitHub 的免费层中,截至 2024 年,网页不能超过 1 GB,页面有每月 100 GB 的软带宽限制。对这个项目来说应该绰绰有余。

Note

有很多不同的平台可以用来调度我们的管道。在我的 ID2223 课程中,学生可以在 Modal 和 GitHub Actions 之间选择。Modal 的免费层很慷慨,开发者体验也很棒,但 Modal 需要信用卡才能访问,而且不能调度 notebooks(只能调度 Python 程序)。还有许多其他提供编排能力的 Serverless 计算平台可以用来运行 Python 程序,包括 Google Cloud Run、Azure Logic Apps、AWS Step Functions、Fly.io、任何托管的 Airflow 平台、Dagster 和 Mage AI。

那么 GitHub Actions 是什么?它是一个持续集成和持续部署(continuous integration and continuous deployment,CI/CD)平台,允许你自动化构建、测试和部署管道。GitHub Actions 通常用于调度测试(单元测试或集成测试)和部署工件。在我们的例子中,我们的特征管道和批推理管道可以被视为部署管道,它们在特征存储中创建特征,并为 GitHub Pages 构建我们的仪表板工件。

为了让你的 GitHub Action 成功运行,你需要把 HOPSWORKS_API_KEY 设置为仓库 secret,这样你的管道才能与 Hopsworks 进行身份验证。

然后你可以继续定义包含 GitHub Actions 的 YAML 文件,它位于 GitHub 仓库的 .github/workflows/air-quality-daily.yml。你可以通过点击"Run workflow"(运行工作流)在你的仓库的 GitHub Actions UI 中运行该工作流。

工作流代码展示了工作流执行的操作。首先,你会注意到这个操作的定时执行已被注释掉。当你成功运行这个 GitHub Action 而无错误后,你可以取消文件开头附近 schedule- cron 行的注释,这个 GitHub Action 就会每天凌晨 6:11 运行。

工作流将执行的步骤如下。首先,工作流将在一个使用最新版 Ubuntu 的容器中运行这些步骤。其次,它会把该 GitHub 仓库中的代码检出到容器中的本地目录,并把当前工作目录改为仓库的根目录。第三,它将安装 Python。第四,它将使用 pip 安装 requirements.txt 文件中的所有 Python 依赖(先把 pip 升级到最新版本)。最后,它将在把 HOPSWORKS_API_KEY 设置为环境变量之后,运行特征管道,然后运行批推理管道。我们的 GitHub Actions 借助 nbconvert 实用程序执行我们的特征管道和批推理 notebooks,它先把 notebook 转换为 Python 程序,然后从第一个单元运行到最后一个单元。设置 HOPSWORKS_API_KEY 环境变量是为了让这些管道能与 Hopsworks 进行身份验证:

`on: workflow_dispatch: #schedule:

- cron: ‘ *11 6 * * **

jobs: test_schedule: runs-on: ubuntu-latest steps: - name: checkout repo content uses: actions/checkout@v4 - name: setup python uses: actions/setup-python@v4 with: python-version: ‘3.10.13’ - name: install python packages run: | python -m pip install –upgrade pip pip install -r requirements.txt - name: execute pipelines env: HOPSWORKS_API_KEY: ${{ *secrets.HOPSWORKS_API_KEY* }} run: | cd notebooks/ch03 jupyter nbconvert –to notebook –execute 2_air_quality_feature_pipeline.ipynb jupyter nbconvert –to notebook –execute 4_air_quality_batch_inference.ipynb`

把仪表板构建为 GitHub Page

我们的 GitHub Action 还包含把批推理管道创建的 PNG 文件提交并推送到我们的 GitHub 仓库的步骤,然后构建并发布包含空气质量预测仪表板(带我们的 PNG 图表)的 GitHub Page。GitHub Action YAML 文件包含一个名为 git-auto-commit-action 的步骤,它把新的 PNG 文件推送到我们的 GitHub 仓库,并重建 GitHub Pages。你不需要更改这段代码:

- name: publish GitHub Pages
        uses: stefanzweifel/git-auto-commit-action@v4 
        [ ... ]

注意,每次 action 运行时,在你的 GitHub 历史中,它都会显示为一次你对仓库的提交。

为了让 git-auto-commit-action 步骤成功运行,你首先必须在你的仓库中启用 GitHub Pages。进入 Settings → Pages → Branch(main → /docs),点击 Save。这将为你的仓库创建 GitHub Page。就这样。一旦你启用了 GitHub Page,并且你的 GitHub Action 每天运行你的工作流,你的仪表板就会每天更新最新的空气质量预报!

使用 LLM 进行函数调用

你现在应该有一个由 ML 驱动的、可用的空气质量预测系统了。但我们想通过添加一个语音激活的 UI 让它更容易使用。为此,我们将使用两个不同的开源 transformer 模型(见图 3-4和仓库中的 5_function_calling.ipynb notebook):

  1. Whisper 把音频转写成文本——用户对着我们的应用说话并提问,模型把用户说的话输出为文本。
  2. 转写后的文本将被输入一个微调过的 Llama 3 8B LLM,它将返回一个函数(从四个可用函数中选出),包括该函数的参数值。
  3. 所选函数将被执行,返回历史空气质量测量值或空气质量预报,该输出将作为提示的一部分,连同你最初的语音提问一起反馈给同一个 Llama 3 8B LLM。
  4. LLM 将返回一个人类能理解的关于空气质量(是否安全或健康)的回答,而不只是 PM 2.5 水平。

流程图展示了使用语音激活系统回答空气质量问题的过程:转写查询、用可用函数查询 LLM、执行函数,并提供用户友好的结果。

原书插图

我们正在使用 RAG(检索增强生成,Retrieval-Augmented Generation)范式构建我们的语音激活 UI,并结合 LLM 函数调用。对于 LLM,用户输入一些文本,称为提示(prompt),LLM 返回一个响应。对于基于聊天的 LLM,比如 OpenAI 的 ChatGPT,响应通常是对话风格的。使用 LLM 进行函数调用时,用户输入一个提示,但现在 LLM 将响应一个 JSON 对象,其中包含要执行的函数(从一组可用函数中选出)以及要传递给该函数的参数。我们将使用一个经过微调、能返回描述函数的 JSON 对象的 LLM。然后我们可以解析 JSON 对象,用它执行我们的预定义函数之一:

  • get_future_data_for_date
  • get_future_data_in_date_range
  • get_historical_air_quality_for_date
  • get_historical_data_in_date_range

也就是说,用户将无法获得关于空气质量的任意问题的答案——只能获得历史读数和空气质量预报。你可以问这样的问题:“上个月的空气质量怎么样?“或"星期二空气质量会怎么样?”

在查询中把函数声明列表传给函数调用 LLM 后,它会尝试用提供的函数之一回答用户查询。LLM 通过分析函数的声明来理解函数的目的。模型实际上不会调用函数。相反,你解析响应,调用模型返回的函数。

这是我们在提示中提供的两个预报函数。另外两个历史函数没有在这里展示,因为它们有类似的定义。注意它们相当冗长,有易于理解的参数名、描述,以及所有参数和返回值的说明:

def get_future_data_for_date \
    (date: str, city_name: str, feature_view, model) -> pd.DataFrame:
    """
    Predicts PM2.5 data for a date and city, given feature view and model.

    Args:
        date (str): The target future date in the format 'YYYY-MM-DD'.
        city_name (str): The name of the city for which the prediction is made.
        feature_view: The feature view used to retrieve batch data.
        model: The machine learning model used for prediction.

    Returns:
        pd.DataFrame: predicted PM2.5 values for each day from target date.

    """

def get_future_data_in_date_range(date_start: str, date_end: str, \
    city_name: str, feature_view, model) -> pd.DataFrame:
    """
    Retrieve data for a specific date range and city from a feature view.

    Args:
        date_start (str): The start date in the format "%Y-%m-%d".
        date_end (str): The end date in the format "%Y-%m-%d".
        city_name (str): The name of the city to retrieve data for.
        feature_view: The feature view object.
        model: The machine learning model used for prediction.

    Returns:
        pd.DataFrame: data for the specified date range and city.
    """

我们为发给我们 LLM 的函数调用查询设计了以下提示模板。首先,我们定义了可用函数,然后包含这些函数的 JSON 表示,包括它们的参数、类型和描述。微调过的 LLM 还应该收到关于选择哪个函数的提示,并被告诉:除非它确信其中一个函数与用户查询匹配,否则不要返回函数:

prompt = f"""<|im_start|>system
You are a helpful assistant with access to the following functions:

get_future_data_for_date
get_future_data_in_date_range
get_historical_air_quality_for_date
get_historical_data_in_date_range

{serialize_function_to_json(get_future_data_for_date)}
{serialize_function_to_json(get_future_data_in_date_range)}
{serialize_function_to_json(get_historical_air_quality_for_date)}
{serialize_function_to_json(get_historical_data_in_date_range)}

You need to choose what function to use and retrieve parameters 
for this function from the user input.
Today is {datetime.date.today().strftime("%A")}, {datetime.date.today()}.
IMPORTANT: If the user query contains 'will', it is very likely that you 
will need to use the get_future_data function.
NOTE: Ignore the Feature View and Model parameters.
NOTE: Dates should be provided in the format YYYY-MM-DD.

To use these functions respond with:
<multiplefunctions>
    <functioncall> {fn} </functioncall>
    <functioncall> {fn} </functioncall>
    ...
</multiplefunctions>

Edge cases you must handle:
- If there are no functions that match the user request, 
you will respond politely that you cannot help.<|im_end|>
<|im_start|>user
{prompt}<|im_end|>
<|im_start|>assistant"""

第二个 LLM 查询的提示可以在源代码仓库中找到。这里不展示,因为它很简单——它包含函数调用的结果、原始用户查询、一些关于空气质量问题的领域知识,以及今天的日期。

5_function_calling.ipynb notebook 需要一个 GPU 才能高效运行。它还有自己的一组需要安装的 Python 依赖:

pip install -r requirements-llm.txt

如果你的笔记本电脑上没有 GPU,你可以免费使用带 T4 GPU 的 Google Colab(不过你需要一个 Google 账户)。你需要取消注释并运行 notebook 中的前两个单元来安装 LLM Python 依赖并下载一些 Python 模块。该 notebook 将 Llama 3 8B 中的权重量化到 4 位,减小了它在内存中的体积,使 LLM 可以在 T4 GPU(有 16 GB RAM)上运行。对我们的系统来说,权重量化似乎不会对 LLM 性能产生负面影响。

还有一个 Streamlit 程序(streamlit_app.py),把同一个 LLM 程序包装在 UI 中。Streamlit 是一个用 Python 以命令式程序构建 UI 的框架。你可以把它托管在免费的无服务器服务上,比如 streamlit.iohuggingface.co

总结与练习

在本章中,我们一起构建了我们的第一个 AI 系统——一个空气质量预测服务。我们把问题分解成总共五个 Python 程序——一个创建和回填特征组的程序、一条下载空气质量读数和天气预报的运营特征管道、一条按需运行的模型训练管道、一条输出空气质量预报图和后报 PNG 文件的批推理管道,以及一个为我们服务提供语音驱动 UI 的 LLM 程序。我们还定义了一个 GitHub Action 工作流(YAML 文件),用于调度特征管道和批推理管道每日运行。这是一大块工作,但现在你有了一个你和你的社区都可以为之自豪的 AI 系统。

以下练习将帮助你学习如何迭代改进你的空气质量预测系统:

  • 给空气质量预测模型添加一个滞后(lagged)PM 2.5 特征。先添加昨天的 PM 2.5 值,然后看看两天或三天前的是否有助于提高模型准确性。
  • 确定添加历史 PM 2.5 值来预测未来 PM 2.5 值存在哪些风险。

1 你可以通过囊性纤维化基金会支持囊性纤维化研究。

2 大文件应该存储在高度可用、可扩展的分布式存储中,比如兼容 S3 的对象存储。这些目前也是存储大文件最便宜的地方。

第二部分 特征存储(Part II. Feature Stores)