1. 从零搭建AI工程体系,为什么我劝你别急着调包
“ai-engineering-from-scratch”这个标题,第一次看到的时候我愣了一下。市面上讲AI的文章铺天盖地,但绝大多数都在教你调API、跑demo、微调个模型就发朋友圈。真正从工程角度,把AI系统当成一个需要设计、需要运维、需要持续迭代的软件项目来对待的内容,少得可怜。
我自己在这个坑里摸爬滚打了几年,带过几个从零到一的AI项目,也接手过别人做了一半跑不起来的烂摊子。说实话,大部分所谓的“AI项目”死掉,不是因为模型不够好,而是因为工程没做好。数据管道是断的,特征存储是散的,模型上线之后没人监控,出了问题只能靠重启。这些问题,调包解决不了。
所以这篇内容,我想认真聊聊“从零构建AI工程体系”这件事。它适合谁看?如果你是一个有后端或数据开发经验、想系统性地理解AI项目该怎么落地的工程师,或者你是一个带团队的技术负责人、正在规划第一个AI产品的基础设施,再或者你是一个数据科学家、发现自己的模型在notebook里跑得挺好但一上线就各种问题——那这篇内容就是写给你的。
我会从整体设计思路开始拆,然后深入到每个核心模块的实操细节,包括数据管道怎么搭、特征怎么管、模型怎么部署、监控怎么做。中间会穿插大量我踩过的坑和实际参数选择的依据。不扯虚的,全是能直接抄作业的东西。
2. 整体架构设计:先想清楚数据怎么流,再想模型怎么跑
2.1 为什么“模型优先”的思路一定会翻车
我见过太多团队一上来就讨论“用哪个模型”“要不要上大模型”“微调还是RAG”,讨论了两周,代码一行没写。这是典型的模型优先思维。正确的顺序应该是反过来的:先想清楚数据从哪来、怎么处理、怎么存储、怎么服务,最后才是模型选型。
原因很简单。在一个AI系统里,模型只是其中一个组件,而且往往是迭代最快的那个组件。今天用BERT,明天可能换成本地部署的小模型,后天可能加一个规则引擎做兜底。如果你的架构是围绕某个具体模型设计的,换模型就意味着重写整个系统。但如果你的架构是围绕数据流设计的,模型就是一个可插拔的模块,换起来成本极低。
我自己的经验是,一个健康的AI工程体系应该分成四层:数据层、特征层、模型层、服务层。每一层之间有清晰的接口,层与层之间通过标准化的数据格式通信。这样你才能做到“模型随便换,管道不用动”。
2.2 四层架构的具体职责划分
数据层负责原始数据的采集、清洗、存储。这里的关键是不可变原始数据原则——原始数据一旦落盘,永远不修改。所有的清洗、转换都在下游做。这样做的好处是,当你发现数据处理逻辑有bug时,可以重新跑一遍管道,而不需要重新采集数据。
特征层是很多人会忽略的一层。它的核心职责是把原始数据转换成模型可以消费的特征,并且保证训练时和推理时的特征计算逻辑完全一致。这一层最典型的坑就是“训练-服务偏差”(training-serving skew),训练时用Python算的特征,推理时用Java算,结果算出来的值不一样,模型效果直接崩掉。
模型层负责模型的训练、评估、版本管理。这里的关键是可复现性——给定同样的数据和代码,必须能训练出同样的模型。这要求你记录每次训练的数据版本、代码版本、超参数、环境依赖。
服务层负责模型的部署和推理。核心要求是低延迟、高可用、可回滚。一个模型上线之后发现效果不好,必须能在分钟级别回滚到上一个版本。
2.3 技术选型的核心考量:别为了时髦买单
在技术选型上,我的原则是:成熟优先,社区活跃度第二,性能第三。AI工程领域的新工具层出不穷,但很多工具的生命周期只有一两年。你不想你的系统上线半年后发现依赖的库已经没人维护了。
具体来说,数据层我用PostgreSQL加对象存储的组合。PostgreSQL存结构化元数据,对象存储存原始文件和大型数据。特征层用Feast或者自己写一套基于Redis加PostgreSQL的方案。模型层用MLflow做实验追踪和模型注册。服务层用FastAPI加Docker加Kubernetes。
这套组合不酷,但极其稳定。我试过用一些更新潮的方案,比如某些流式特征平台,结果发现调试成本极高,出了问题连日志都找不到。对于大多数团队来说,能跑通、能调试、能维护比什么都重要。
提示:如果你团队规模小于5人,不要自建特征平台。直接用PostgreSQL加定时任务算特征,简单可靠。等你的特征数量超过200个、或者需要实时特征的时候再考虑引入专门的平台。
3. 数据管道搭建:从原始数据到可用数据集的关键步骤
3.1 数据采集与落盘的工程细节
数据采集听起来简单,但实际操作中有大量细节。首先是幂等性——同一条数据重复采集不能产生重复记录。我的做法是在采集端生成一个基于内容哈希的唯一ID,写入时用upsert语义。这样即使采集任务重跑,也不会污染数据。
其次是时间戳管理。每条数据必须有两个时间戳:事件发生时间(event_time)和入库时间(ingest_time)。这两个时间戳在后续的特征计算和问题排查中至关重要。我踩过的坑是早期只记录了入库时间,结果做时间窗口特征时发现数据延迟导致特征穿越(leakage),模型离线评估AUC 0.95,上线后掉到0.6。
第三是数据分区。按天分区是最常见的做法,但如果你的数据量很大,建议按小时分区。分区的好处是查询时可以裁剪掉不需要的数据,而且删除过期数据时可以直接drop分区,比delete快几个数量级。
-- 按天分区的数据表结构示例 CREATE TABLE raw_events ( event_id TEXT PRIMARY KEY, event_time TIMESTAMPTZ NOT NULL, ingest_time TIMESTAMPTZ DEFAULT NOW(), user_id TEXT NOT NULL, event_type TEXT NOT NULL, payload JSONB ) PARTITION BY RANGE (event_time); -- 创建每日分区 CREATE TABLE raw_events_20240101 PARTITION OF raw_events FOR VALUES FROM ('2024-01-01') TO ('2024-01-02');3.2 数据清洗的标准化流程
数据清洗不是“把脏数据删掉”这么简单。我的标准流程是四步:检测、标记、修复、验证。
检测阶段,我会对每个字段计算一组统计指标:空值率、唯一值数量、数值分布的分位数、字符串长度分布。这些指标和上一周期的指标对比,如果偏差超过阈值就告警。比如某个字段的空值率从2%突然跳到30%,大概率是上游出了问题。
标记阶段,不是直接删除异常数据,而是给每条记录打上质量标签。比如quality_flag字段,值为ok、missing_required、out_of_range、duplicate等。这样后续处理时可以根据标签决定怎么处理,而不是一刀切。
修复阶段,对于缺失值,根据业务逻辑决定是填充默认值、用统计量填充、还是丢弃。对于异常值,我倾向于winsorize(缩尾处理)而不是直接删除。比如用户年龄出现200岁,缩到99岁比删掉这条记录更合理,因为这条记录的其他字段可能是有价值的。
验证阶段,清洗后的数据要重新跑一遍统计指标,确认清洗逻辑达到了预期效果。这一步经常被跳过,但它是保证数据质量的关键。
3.3 数据版本管理:别让“数据变了”成为玄学
数据版本管理是AI工程和传统软件工程最大的区别之一。代码有Git管理,但数据呢?很多团队的数据是“活的”——每天都在变,导致上周跑出好结果的实验这周复现不了。
我的做法是快照加增量。每天凌晨对核心数据表做一次快照,快照存储为Parquet格式,放在对象存储上,路径包含日期和版本号。训练时指定数据版本,比如dataset_v20240101。这样任何时候都能复现当时的训练数据。
对于增量数据,记录每次变更的元信息:变更时间、变更行数、变更原因。这些元信息存在一个专门的元数据表里,方便追溯。
# 数据快照的简单实现 import pandas as pd from datetime import date def create_snapshot(table_name: str, snapshot_date: str): query = f""" SELECT * FROM {table_name} WHERE ingest_time < '{snapshot_date}' """ df = pd.read_sql(query, con=engine) # 计算数据指纹 fingerprint = hashlib.md5( pd.util.hash_pandas_object(df).values.tobytes() ).hexdigest() # 存储快照 path = f"s3://data-snapshots/{table_name}/{snapshot_date}/data.parquet" df.to_parquet(path) # 记录元数据 record_metadata(table_name, snapshot_date, len(df), fingerprint) return fingerprint注意:快照不是备份。快照的目的是复现实验,备份的目的是灾难恢复。两者策略不同,不要混为一谈。
4. 特征工程与特征存储:训练和推理一致性的保障
4.1 特征定义的标准模板
特征工程最怕的是“口口相传”——某个人在notebook里写了一段特征计算代码,另一个人复制粘贴到另一个notebook里,改了几个参数,然后两个版本的特征定义就不一致了。解决这个问题的唯一办法是特征定义代码化、模板化。
我给每个特征定义一个标准的Python类,包含以下要素:特征名称、数据类型、计算逻辑、依赖的原始字段、时间窗口、默认值、负责人。这个类同时被训练管道和推理服务引用,保证逻辑一致。
from dataclasses import dataclass from typing import List, Optional import pandas as pd @dataclass class FeatureDefinition: name: str dtype: str dependencies: List[str] window: Optional[str] = None default_value: any = None owner: str = "" def compute(self, df: pd.DataFrame) -> pd.Series: raise NotImplementedError class UserPurchaseCount7d(FeatureDefinition): def __init__(self): super().__init__( name="user_purchase_count_7d", dtype="int64", dependencies=["user_id", "event_time", "event_type"], window="7d", default_value=0, owner="data-team" ) def compute(self, df: pd.DataFrame) -> pd.Series: cutoff = df["event_time"].max() - pd.Timedelta(days=7) mask = (df["event_type"] == "purchase") & (df["event_time"] >= cutoff) return df[mask].groupby("user_id").size()这个模板的好处是,特征的计算逻辑只有一份,训练时批量计算,推理时单条计算,但走的是同一套代码。我实测下来,这套方案把训练-服务偏差导致的问题减少了90%以上。
4.2 离线特征与在线特征的同步策略
离线特征用于训练,在线特征用于推理。两者的存储介质不同:离线特征存在数据仓库或Parquet文件里,在线特征存在Redis或内存数据库中。关键问题是:怎么保证两者一致?
我的策略是离线优先,在线派生。所有特征先离线计算好,写入离线存储。然后通过一个同步任务,把最新版本的特征推送到在线存储。在线存储只保留每个实体的最新特征值,不保留历史。
同步任务的设计要点:第一,必须幂等,重复执行不会产生错误;第二,必须有版本号,在线特征带一个feature_version字段,推理时可以检查版本是否匹配;第三,必须有回滚机制,同步出错时能快速恢复到上一个版本。
# 特征同步任务的核心逻辑 def sync_features_to_online(feature_names: List[str], batch_size: int = 1000): for feature_name in feature_names: # 从离线存储读取最新特征 df = read_offline_features(feature_name) # 分批写入在线存储 for i in range(0, len(df), batch_size): batch = df.iloc[i:i+batch_size] pipe = redis_client.pipeline() for _, row in batch.iterrows(): key = f"feature:{feature_name}:{row['entity_id']}" value = { "value": row["feature_value"], "version": row["feature_version"], "updated_at": row["updated_at"] } pipe.set(key, json.dumps(value)) pipe.execute() # 记录同步元数据 record_sync_metadata(feature_name, len(df), datetime.now())4.3 特征监控:及时发现特征漂移
特征漂移是模型效果下降的头号原因。用户行为变了、市场环境变了、上游数据源变了,都会导致特征分布发生变化。如果不监控,你只能等到业务指标下降才发现问题,那时候已经晚了。
我监控三类指标:统计指标(均值、方差、分位数)、分布指标(PSI、KL散度)、缺失率。统计指标每天计算一次,和过去30天的均值对比,偏差超过2个标准差就告警。分布指标每周计算一次,PSI超过0.2就告警。缺失率实时监控,超过阈值立即告警。
import numpy as np from scipy import stats def calculate_psi(expected: np.ndarray, actual: np.ndarray, buckets: int = 10) -> float: """计算群体稳定性指数(PSI)""" breakpoints = np.percentile(expected, np.linspace(0, 100, buckets + 1)) breakpoints[0] = -np.inf breakpoints[-1] = np.inf expected_counts = np.histogram(expected, bins=breakpoints)[0] / len(expected) actual_counts = np.histogram(actual, bins=breakpoints)[0] / len(actual) # 避免除零 expected_counts = np.where(expected_counts == 0, 0.0001, expected_counts) actual_counts = np.where(actual_counts == 0, 0.0001, actual_counts) psi = np.sum((actual_counts - expected_counts) * np.log(actual_counts / expected_counts)) return psi提示:PSI的阈值不是绝对的。0.1以下说明分布稳定,0.1到0.2说明有轻微变化,0.2以上说明分布显著变化。但具体阈值要根据业务场景调整。金融风控场景可能0.1就要告警,推荐系统场景0.25才需要关注。
5. 模型训练与部署:从实验到上线的完整链路
5.1 实验追踪:别让好结果找不到
我见过最可惜的事情是:一个数据科学家跑出了一个效果很好的模型,但两周后想复现时发现找不到当时的代码、参数和数据版本。实验追踪不是可选项,是必选项。
MLflow是我用得最顺手的工具。每次训练自动记录:代码的Git commit hash、数据版本、超参数、评估指标、模型文件、环境依赖。这些信息存在一个中心化的数据库里,任何人都可以查询。
import mlflow import mlflow.sklearn def train_model(params: dict, data_version: str): mlflow.set_experiment("my-ai-project") with mlflow.start_run(): # 记录数据版本 mlflow.log_param("data_version", data_version) # 记录超参数 mlflow.log_params(params) # 记录Git commit commit_hash = subprocess.check_output( ["git", "rev-parse", "HEAD"] ).decode().strip() mlflow.log_param("git_commit", commit_hash) # 训练模型 model = train(data_version, params) # 记录指标 metrics = evaluate(model) mlflow.log_metrics(metrics) # 记录模型 mlflow.sklearn.log_model(model, "model") return model5.2 模型注册与版本管理
模型注册表是模型从实验走向生产的桥梁。每个模型版本有四个状态:staging(测试中)、production(生产中)、archived(已归档)、none(未指定)。状态转换需要审批,不能随便改。
我的做法是:训练完成后模型自动进入staging状态,在staging环境跑一周的A/B测试,效果达标后手动提升到production。提升时记录提升人、提升时间、提升原因。回滚时同样记录。
模型版本号我用语义化版本:major.minor.patch。major版本表示模型架构变化,minor版本表示特征或超参数变化,patch版本表示训练数据更新。这样从版本号就能大致判断变更的影响范围。
5.3 部署模式选择:批处理、实时、还是流式
部署模式的选择取决于业务需求。我总结了一个简单的决策表:
| 场景 | 推荐模式 | 延迟要求 | 实现复杂度 |
|---|---|---|---|
| 离线报表 | 批处理 | 小时级 | 低 |
| 用户画像更新 | 批处理+缓存 | 分钟级 | 中 |
| 实时推荐 | 实时推理 | 毫秒级 | 高 |
| 风控决策 | 实时推理 | 毫秒级 | 高 |
| 内容审核 | 流式推理 | 秒级 | 中高 |
批处理最简单,用Airflow或cron定时跑就行。实时推理需要模型服务常驻内存,用FastAPI加Uvicorn部署。流式推理介于两者之间,用Kafka加消费者组实现。
我建议从批处理开始。很多团队一上来就搞实时推理,结果发现业务根本不需要那么低的延迟,白白增加了系统复杂度和维护成本。等业务确实需要实时了再迁移,迁移成本没有想象中那么高。
5.4 模型服务的性能优化
模型服务上线后,性能是第一个要关注的问题。我遇到过模型推理延迟从50ms涨到500ms的情况,排查后发现是特征获取环节出了问题——每次推理都要查一次数据库,数据库连接池被打满了。
优化手段有几个:第一,特征预取。对于每个请求,一次性批量获取所有需要的特征,而不是逐个获取。第二,连接池。数据库连接、Redis连接都要用连接池,避免频繁创建销毁。第三,批处理推理。如果延迟要求允许,把多个请求攒成一批一起推理,GPU利用率能提升好几倍。
from fastapi import FastAPI from contextlib import asynccontextmanager import asyncpg import aioredis @asynccontextmanager async def lifespan(app: FastAPI): # 启动时创建连接池 app.state.db_pool = await asyncpg.create_pool( dsn="postgresql://...", min_size=10, max_size=50 ) app.state.redis = await aioredis.create_redis_pool( "redis://...", minsize=10, maxsize=50 ) yield # 关闭时释放连接 await app.state.db_pool.close() app.state.redis.close() app = FastAPI(lifespan=lifespan) @app.post("/predict") async def predict(request: PredictRequest): # 批量获取特征 features = await get_features_batch( app.state.redis, request.entity_ids ) # 批量推理 predictions = model.predict_batch(features) return {"predictions": predictions.tolist()}注意:连接池的大小不是越大越好。PostgreSQL默认最大连接数是100,如果你的服务有多个实例,每个实例的连接池大小要除以实例数。否则会出现连接数超限的问题。
6. 监控与运维:让AI系统稳定运行的关键
6.1 模型效果监控的指标体系
模型上线不是终点,而是起点。你需要持续监控模型的效果。我监控的指标分三层:业务指标、模型指标、系统指标。
业务指标是最重要的,比如点击率、转化率、GMV。这些指标直接反映模型对业务的影响。模型指标包括AUC、准确率、召回率等,用于判断模型本身是否退化。系统指标包括延迟、吞吐量、错误率,用于判断服务是否健康。
三层指标的关系是:系统指标异常会导致模型指标异常,模型指标异常会导致业务指标异常。所以排查问题时,从系统指标开始查,逐层往上。
6.2 告警策略:别让告警疲劳毁掉运维
告警太多等于没有告警。我见过一个系统每天发几百条告警,运维人员直接屏蔽了告警群,结果真正出问题时没人知道。
我的告警策略是分级加聚合。P0告警(业务指标下降超过10%)立即电话通知。P1告警(模型指标下降超过5%)发到告警群,30分钟内响应。P2告警(系统指标异常但未影响业务)发到告警群,当天处理。P3告警(趋势性变化)每天汇总一次。
聚合的意思是,同一类型的告警在5分钟内只发一条,避免刷屏。比如特征缺失率告警,如果10个特征同时缺失,只发一条汇总告警,而不是10条。
6.3 常见故障排查速查表
| 现象 | 可能原因 | 排查方法 | 解决方案 |
|---|---|---|---|
| 推理延迟突增 | 特征获取慢 | 查看特征存储的响应时间 | 加缓存、优化查询 |
| 模型效果下降 | 特征漂移 | 计算PSI、对比特征分布 | 重新训练模型 |
| 服务报错率上升 | 依赖服务故障 | 查看依赖服务的健康状态 | 降级、熔断 |
| 预测结果异常 | 输入数据格式变化 | 检查输入数据的schema | 修复上游数据管道 |
| 内存溢出 | 批处理大小过大 | 查看内存使用曲线 | 减小批处理大小 |
| GPU利用率低 | 批处理大小过小 | 查看GPU利用率 | 增大批处理大小 |
这张表是我自己整理的,每次遇到问题先查表,80%的情况能直接定位到原因。剩下的20%需要深入排查,但至少有了方向。
6.4 模型重训练的触发机制
模型不是训练一次就一劳永逸的。什么时候需要重训练?我设定了三个触发条件:定时触发(每周一次)、效果触发(业务指标下降超过阈值)、数据触发(特征漂移超过阈值)。
定时触发是兜底,保证模型至少每周更新一次。效果触发是核心,业务指标下降说明模型已经不适应当前数据了。数据触发是预警,特征漂移说明数据分布变了,模型可能很快就不适用了。
重训练不是全自动的。我的流程是:触发重训练后,自动跑一遍训练管道,生成新模型,在staging环境评估。评估通过后通知人工审核,人工确认后才上线。这样既保证了及时性,又避免了自动上线带来的风险。
7. 我踩过的坑和给你的建议
7.1 那些年我踩过的数据坑
最大的坑是数据穿越。早期做用户流失预测,我用了一个特征叫“用户最近一次登录距今天数”。训练时用全量数据算这个特征,结果模型学到了“登录天数少的人会流失”这个显而易见的规律,离线AUC 0.92。上线后发现效果很差,因为推理时这个特征是用当前时间算的,而训练时用的是数据截止时间,两者不一致。
修复方法是:所有时间窗口特征必须基于一个统一的参考时间点计算。训练时参考时间点是样本的标签时间,推理时参考时间点是当前时间。这个逻辑必须写在特征定义里,不能靠人工保证。
第二个坑是数据延迟。上游数据源延迟了6小时,但特征计算任务按小时调度,导致最近6小时的特征全是默认值。模型在推理时拿到这些默认值,预测结果自然不准。解决方案是加一个数据新鲜度检查,如果数据延迟超过阈值,特征计算任务跳过,推理服务使用缓存的特征值。
7.2 模型部署的常见陷阱
模型部署最大的陷阱是环境不一致。训练时用Python 3.9加PyTorch 1.12,推理时用Python 3.10加PyTorch 2.0,模型加载直接报错。解决方案是用Docker把训练环境和推理环境统一起来,用同一个基础镜像。
第二个陷阱是模型文件过大。一个BERT模型动辄几百MB,加载一次要几十秒。如果服务重启,这段时间内所有请求都会失败。解决方案是模型预热——服务启动时先加载模型跑一次推理,确认模型可用后再接收流量。
第三个陷阱是版本回滚困难。新模型上线后发现效果不好,想回滚到旧版本,结果发现旧版本的模型文件已经被覆盖了。解决方案是模型文件按版本号存储,永不删除。存储成本很低,但回滚时的价值极高。
7.3 给刚入门的团队的建议
如果你刚开始做AI工程,我的建议是:先跑通一个最小闭环,再逐步优化。最小闭环包括:数据采集、特征计算、模型训练、模型部署、效果监控。每个环节先用最简单的方案实现,跑通之后再考虑优化。
不要一开始就追求实时推理、特征平台、自动重训练这些高级功能。这些功能在业务量小的时候不仅没用,还会增加系统复杂度和维护成本。等业务量上来了,痛点自然会出现,到时候再针对性解决。
另外,文档和监控比代码更重要。代码写得好但没人知道怎么用,等于没写。监控做得好但代码写得烂,至少系统还能跑。我见过太多团队把精力全花在代码上,结果系统上线后没人知道怎么排查问题,一出故障就抓瞎。
最后分享一个小技巧:每次上线新模型时,保留1%的流量走旧模型,作为对照组。这样即使新模型有问题,你也能通过对比及时发现。这个策略叫“影子流量”,实现成本很低,但价值极高。我靠这个策略避免了好几次重大事故。