资讯动态

从零搭建AI工程体系:数据流、推理服务与监控闭环实战

发布时间:2026/9/30 12:07:59 来源:尧图企业网站定制
1. 从零搭建AI工程体系为什么我劝你别一上来就调包这两年“AI工程”这个词被聊烂了。打开任何一个技术社区满屏都是“三行代码调用大模型”“十分钟搭建RAG”“零基础微调你的第一个模型”。看着确实爽但我带过几个新人之后发现一个共性问题能跑通Demo的人一大把能扛住线上流量、能把成本压下来、能定位到是数据问题还是模型问题的人少得可怜。ai-engineering-from-scratch这个标题我理解成两层意思。第一层是“从零开始学AI工程”第二层是“从零手搓AI工程链路”。这两层其实是同一件事的两面——你不亲手把一条链路搭一遍就永远不知道那些封装好的框架在背后替你做了什么出了问题也就无从下手。我打算把这条链路完整拆一遍从数据准备、特征处理、模型训练、推理服务到监控和迭代。不依赖任何“一键式”平台核心环节尽量用最朴素的工具实现让你看清楚每一层到底在干什么。适合谁看如果你已经会写Python、懂一点机器学习基础但一到工程落地就发怵那这篇就是写给你的。如果你是完全的新手也能看懂大框架只是细节部分需要你边看边动手补。先说结论AI工程的核心不是模型是数据流和反馈闭环。模型只是这条流水线上的一个零件而且是可以随时替换的零件。把这句话记住后面所有内容都围绕它展开。2. 整体架构设计一条能落地的AI工程链路长什么样2.1 先想清楚为什么不能直接调API了事很多人会问现在大模型API这么便宜我直接调不就完了为什么要自己搭链路这个问题问得好答案取决于你的场景。如果你的需求是“把一段文字润色一下”“做个简单的问答”那确实调API就够了自己搭纯属浪费时间。但如果你面对的是下面这些情况就必须自己掌控链路数据不能出内网金融、医疗、政企类项目数据合规是硬门槛第三方API直接出局。成本随调用量线性增长当你的日调用量到百万级API费用会变成一笔非常可观的支出自建推理服务的边际成本会低很多。需要深度定制通用模型在你的垂直领域表现拉胯必须做微调或者领域适配。延迟要求苛刻API的网络往返加上排队很难稳定做到百毫秒以内自建服务可以做到极致优化。我个人的判断标准很简单先算一笔账再决定要不要自建。下面这张表是我常用的决策参考你可以直接套。场景特征推荐方案理由日调用量 1万数据不敏感直接调API省时省力成本可控日调用量 1万~50万数据不敏感API 缓存 路由平衡成本与开发效率日调用量 50万或数据敏感自建推理服务边际成本低数据可控需要领域深度定制自建 微调通用模型无法满足延迟要求 100ms自建 量化 推理优化API无法保证稳定延迟这张表不是绝对的但能帮你快速定位自己该走哪条路。我见过太多团队一上来就自建结果发现调用量根本撑不起运维成本最后又退回API。先跑通业务再优化成本这个顺序不能反。2.2 分层设计把系统切成五块每块独立演进一条完整的AI工程链路我习惯切成五层。这样切的好处是每一层可以独立替换、独立扩容、独立监控出了问题能快速定位到是哪一层的锅。第一层数据接入层。负责从各种源头把数据捞进来——数据库、日志文件、消息队列、第三方接口。这一层的核心任务是统一数据格式不管上游是什么样进来之后都变成标准的结构化记录。第二层数据处理层。负责清洗、去重、标注、特征提取。这一层是最脏最累的但也是决定模型上限的关键。业内有个说法数据和特征决定了模型的上限算法只是逼近这个上限。我深以为然。第三层模型训练层。负责训练、调参、评估、版本管理。这一层反而相对标准化因为工具链成熟照着最佳实践走就行。第四层推理服务层。负责把训练好的模型部署成可调用的服务处理并发、做批处理、做缓存、做降级。这一层直接面向业务对稳定性和延迟要求最高。第五层监控反馈层。负责监控线上表现、收集bad case、触发再训练。这一层最容易被忽略但它是整个系统能否持续进化的关键。这五层之间通过明确定义的接口通信比如数据处理层输出的是标准格式的训练样本推理服务层输入的是标准格式的请求。接口定好了每层内部怎么实现都可以换。2.3 技术选型够用就好别追新选型这块我踩过不少坑。早年追新什么火用什么结果半年后社区不维护了迁移成本高得吓人。后来学乖了选型的核心原则是成熟、活跃、可替换。具体到每一层我的常用组合是这样的数据接入Python Pandas 处理中小规模Spark 处理大规模。消息队列用 Kafka 或者 RabbitMQ看团队熟悉哪个。数据处理Pandas NumPy 做原型PySpark 做规模化。特征存储可以考虑 Feast但小团队用文件系统加版本号也够。模型训练PyTorch 是首选生态最全。实验管理用 MLflow 或者 Weights Biases小团队用 MLflow 自建省钱。推理服务FastAPI 做接口ONNX Runtime 或者 TensorRT 做加速。并发量大了上 Triton Inference Server。监控反馈Prometheus Grafana 做指标监控日志用 ELK 或者 Loki。这套组合不是最优的但胜在每一块都有大量文档和社区支持出了问题能搜到答案。新手最容易犯的错是选了一堆小众工具结果卡在一个冷门bug上三天出不来。提示选型时优先考虑“如果这个工具明天停止维护我换掉它的成本有多大”。可替换性比性能更重要尤其是在早期。3. 数据层AI工程的地基90%的问题出在这里3.1 数据接入别小看“把数据捞进来”这件事数据接入听起来简单实际做起来坑很多。我总结下来主要是三个问题格式不统一、时序不对齐、增量难处理。格式不统一好理解上游有JSON、有CSV、有数据库表字段名还各不相同。我的做法是定义一个中间表示层所有数据进来先转成这个统一格式再往下游走。这个中间格式用Parquet最合适列式存储、压缩率高、读取快而且Pandas和Spark都原生支持。时序不对齐是个隐蔽的坑。比如你要预测用户行为特征来自多个系统有的系统是实时写入有的系统是T1批量同步。如果你不处理时间戳对齐训练出来的模型在线上会表现得很奇怪——因为线上推理时拿到的特征和训练时的时间关系不一致。解决办法是给每条记录打上事件时间戳所有特征都按事件时间对齐而不是按写入时间。增量处理是另一个难点。全量重跑在数据量小的时候没问题数据量一大就扛不住。我的做法是用水位线watermark机制记录上次处理到哪个时间点每次只处理新增部分。但要注意处理迟到数据——有些数据会晚到如果直接跳过就会丢。通常我会保留一个回看窗口比如每次多处理最近24小时的数据做幂等写入。# 一个简化的增量接入示例 import pandas as pd from datetime import datetime, timedelta def incremental_ingest(source_path, watermark_file, lookback_hours24): # 读取上次水位线 last_watermark pd.read_json(watermark_file).iloc[0][watermark] # 回看窗口处理迟到数据 start_time last_watermark - timedelta(hourslookback_hours) # 只读取新增数据 df pd.read_parquet(source_path, filters[(event_time, , start_time)]) # 幂等写入按主键去重 df df.drop_duplicates(subset[record_id], keeplast) # 更新水位线 new_watermark df[event_time].max() pd.DataFrame({watermark: [new_watermark]}).to_json(watermark_file) return df这段代码的核心是幂等——同样的数据重复处理不会产生副作用。这是增量处理的生命线因为分布式环境下重复消费几乎不可避免。3.2 数据清洗脏数据不处理模型再牛也白搭数据清洗这块我有一套固定的检查清单每次新数据集进来都过一遍缺失值是随机缺失还是系统性缺失随机缺失可以填充系统性缺失往往意味着数据采集有问题得先修采集。异常值是真实异常还是录入错误用分位数和业务规则双重判断。重复值完全重复还是近似重复近似重复要用相似度算法处理。标签噪声标注数据里有多少错标抽样人工复核估算噪声率。分布偏移训练集和验证集分布是否一致用KS检验或者对抗验证检测。这里重点说标签噪声因为它最容易被忽略但危害最大。我做过一个实验在同一个数据集上把5%的标签随机翻转模型准确率从92%掉到78%。而如果把这5%的噪声清理掉准确率能回到90%以上。清理标签噪声的收益往往比换模型架构大得多。清理标签噪声的常用方法有两种。一种是置信学习Confident Learning用交叉验证的预测概率来估计哪些样本可能标错了。另一种是人工复核对模型预测和标签不一致的样本重点检查。小数据集用后者大数据集用前者。3.3 特征工程从原始数据到模型能吃的格式特征工程是数据层最考验功力的地方。我的经验是先把业务逻辑理清楚再谈技术实现。很多特征不是算不出来而是你根本没想到要算。举个例子做用户流失预测。原始数据只有用户ID、注册时间、最近登录时间。直接把这些喂给模型效果一般。但如果你能算出“最近7天登录次数”“登录间隔的方差”“活跃天数的占比”这些衍生特征效果会明显提升。这些特征背后是业务理解——流失的本质是用户行为模式的改变而不是某个时间点的状态。特征工程的技术实现上我推荐用**特征管道Feature Pipeline**的方式把每个特征的計算逻辑封装成独立的函数然后用管道串起来。这样做的好处是可测试、可复用、可版本管理。# 特征管道示例 from sklearn.base import BaseEstimator, TransformerMixin class LoginFrequencyFeature(BaseEstimator, TransformerMixin): def __init__(self, window_days7): self.window_days window_days def fit(self, X, yNone): return self def transform(self, X): # X 包含 user_id, login_time cutoff X[login_time].max() - pd.Timedelta(daysself.window_days) recent X[X[login_time] cutoff] freq recent.groupby(user_id).size().rename(login_freq_7d) return freq class FeaturePipeline: def __init__(self, features): self.features features def fit_transform(self, X): result X.copy() for feat in self.features: result result.join(feat.fit_transform(X), howleft) return result这种写法的好处是每个特征独立可测出问题能快速定位。而且特征定义和计算逻辑在一起不会出现“文档写的是A代码算的是B”这种破事。注意特征工程一定要做训练/推理一致性校验。我见过太多案例训练时用Pandas算特征推理时用Java重写一遍结果两边逻辑有细微差异线上效果直接崩盘。解决办法是特征计算逻辑只写一遍训练和推理共用同一套代码。4. 模型层训练、评估与版本管理4.1 训练流程把实验变成可复现的工程模型训练最容易犯的错是实验不可复现。今天调了个参数效果很好明天想再跑一遍结果忘了当时改了什么或者数据变了再也复现不出来。这在工程上是灾难。我的做法是三固定固定数据版本、固定代码版本、固定环境版本。数据版本每次训练用的数据集打上版本号存到对象存储或者特征存储里。用DVC或者LakeFS管理都行。代码版本Git commit hash 记录到实验元数据里。环境版本用Docker镜像锁定依赖版本requirements.txt 不够因为系统库版本也会影响。这三样固定了任何一次实验都能精确复现。MLflow 这类工具能帮你自动记录这些信息但核心是你要有意识地去记录。训练流程本身我习惯分成四个阶段数据加载、模型构建、训练循环、评估保存。每个阶段都有明确的输入输出方便单独调试。# 训练流程骨架 import mlflow import torch def train(config): with mlflow.start_run(): # 记录配置 mlflow.log_params(config) # 1. 数据加载 train_loader, val_loader build_dataloaders(config) # 2. 模型构建 model build_model(config) optimizer build_optimizer(model, config) # 3. 训练循环 for epoch in range(config[epochs]): train_loss train_one_epoch(model, train_loader, optimizer) val_loss, val_metrics evaluate(model, val_loader) mlflow.log_metrics({ train_loss: train_loss, val_loss: val_loss, **val_metrics }, stepepoch) # 早停检查 if early_stopping(val_loss): break # 4. 保存模型 mlflow.pytorch.log_model(model, model) return model这个骨架看起来简单但每个环节都有讲究。比如数据加载阶段shuffle 的顺序要固定随机种子否则每次训练的数据顺序不一样结果会有波动。再比如评估阶段验证集不能参与任何训练决策包括早停。早停应该用独立的验证集或者用交叉验证。4.2 评估体系别只看准确率评估模型不能只看一个指标。我通常从四个维度评估区分度模型能不能把正负样本分开AUC、KS 是常用指标。校准度模型输出的概率准不准用校准曲线和Brier分数衡量。稳定性模型在不同数据切片上表现是否一致做分组评估。业务指标模型上线后对业务KPI的影响这个最重要但往往被忽略。校准度是很多人忽略的维度。举个例子模型说“这个用户有80%的概率会流失”如果实际流失率只有50%那这个概率就是不准的。在需要根据概率做决策的场景比如发优惠券的阈值校准度比区分度更重要。校准的方法有 Platt Scaling 和 Isotonic Regression前者适合小数据集后者适合大数据集。我一般先用 Isotonic如果数据量不够再退回 Platt。分组评估也很关键。整体AUC 0.85 看着不错但如果拆开看新用户AUC 0.9老用户AUC 0.6那这个模型对老用户基本没用。分组维度可以是用户类型、地域、时间段看业务关心什么。4.3 版本管理模型也要有“身份证”模型版本管理不是简单地存个文件。一个完整的模型版本应该包含模型权重训练好的参数。模型结构网络定义或者配置文件。预处理逻辑特征工程代码必须和训练时一致。后处理逻辑输出转换代码。依赖环境Docker镜像或者conda环境文件。评估报告在哪些数据上评估的指标是多少。训练元数据数据版本、代码版本、超参数。这些东西打包在一起才是一个可部署、可回滚的模型版本。我见过团队只存了权重文件结果要回滚的时候发现预处理代码找不到了只能重新训练白白浪费几天。版本命名我推荐用语义化版本 时间戳比如v1.2.0-20240115。语义化版本表示功能变更时间戳保证唯一性。回滚的时候按时间戳找最近的稳定版本。5. 推理服务层从模型文件到线上服务5.1 服务框架FastAPI 起步按需升级推理服务的框架选择我的建议是从 FastAPI 开始。它足够简单性能也不差配合 Uvicorn 跑异步单机扛几百QPS没问题。等真的扛不住了再考虑上 Triton 或者自己写 C 服务。FastAPI 服务的基本结构是这样的from fastapi import FastAPI from pydantic import BaseModel import torch app FastAPI() class PredictRequest(BaseModel): features: list[float] class PredictResponse(BaseModel): score: float version: str # 启动时加载模型只加载一次 model None MODEL_VERSION v1.0.0 app.on_event(startup) def load_model(): global model model torch.jit.load(model.pt) model.eval() app.post(/predict, response_modelPredictResponse) async def predict(req: PredictRequest): with torch.no_grad(): tensor torch.tensor([req.features]) score model(tensor).item() return PredictResponse(scorescore, versionMODEL_VERSION)这里有几个关键点。模型在启动时加载不要每次请求都加载否则延迟会高得离谱。用torch.no_grad()关闭梯度计算推理时不需要梯度关掉能省内存和计算。返回版本号方便排查问题时确认是哪个版本的服务。5.2 性能优化延迟和吞吐的平衡术推理服务的性能优化核心是批处理batching。单个请求推理一次GPU利用率很低。把多个请求攒成一批一起推理吞吐能提升几倍甚至几十倍。但批处理会引入延迟——你得等一批攒够才能推理。所以要在延迟和吞吐之间找平衡。我的做法是设置最大等待时间和最大批大小两个参数谁先到就触发推理。import asyncio from collections import deque class BatchProcessor: def __init__(self, model, max_batch32, max_wait0.01): self.model model self.max_batch max_batch self.max_wait max_wait self.queue deque() self.lock asyncio.Lock() async def predict(self, features): future asyncio.Future() async with self.lock: self.queue.append((features, future)) if len(self.queue) self.max_batch: await self._process() else: asyncio.create_task(self._delayed_process()) return await future async def _delayed_process(self): await asyncio.sleep(self.max_wait) async with self.lock: if self.queue: await self._process() async def _process(self): batch list(self.queue) self.queue.clear() features [f for f, _ in batch] results self.model(features) for (_, future), result in zip(batch, results): future.set_result(result)这段代码实现了一个简单的动态批处理。max_batch32表示最多攒32个请求max_wait0.01表示最多等10毫秒。实际参数要根据你的延迟要求和流量特征调。除了批处理还有几个优化手段模型量化FP16或INT8精度损失小但速度提升明显、算子融合用ONNX Runtime或TensorRT自动做、缓存对相同输入直接返回缓存结果。这些手段可以叠加使用但要注意量化可能带来精度损失上线前要评估。5.3 降级与容错线上服务不能挂推理服务最怕的是雪崩。一个请求超时导致线程池占满后续请求全部排队最后整个服务不可用。防止雪崩的核心是快速失败 降级。我的做法是给每个请求设置超时时间超时直接返回默认值或者走降级逻辑。降级逻辑可以是返回缓存结果、返回规则引擎的结果、或者返回一个保守的默认值。import asyncio from fastapi import HTTPException app.post(/predict) async def predict(req: PredictRequest): try: result await asyncio.wait_for( model_inference(req.features), timeout0.5 # 500ms 超时 ) return result except asyncio.TimeoutError: # 降级返回缓存或默认值 fallback get_fallback_result(req.features) return fallback除了超时降级还要做限流。用令牌桶或者漏桶算法限制每秒处理的请求数。超过限制的请求直接拒绝返回429状态码。这样能保护后端不被压垮。健康检查也必不可少。Kubernetes 会定期调健康检查接口如果服务不健康就自动重启或者摘除流量。健康检查要检查关键依赖比如模型是否加载成功、数据库是否连通。6. 监控与迭代让系统自己进化6.1 监控指标看什么怎么看线上监控我分四类指标系统指标CPU、内存、GPU利用率、网络IO。这些用 Prometheus Grafana 采集展示。服务指标QPS、延迟分布P50/P95/P99、错误率。这些在服务代码里埋点。模型指标预测分布、特征分布、置信度分布。这些用来检测数据漂移。业务指标转化率、点击率、留存率。这些是最终衡量模型价值的标准。延迟分布比平均延迟重要得多。平均延迟100ms看着不错但如果P99是5秒那1%的用户体验极差。我通常关注P95和P99这两个指标能反映长尾情况。数据漂移检测是模型监控的核心。线上数据的分布会随时间变化如果训练时的分布和线上分布差异太大模型效果会下降。检测方法有PSIPopulation Stability Index和KS检验。PSI大于0.2通常认为有显著漂移需要触发再训练。import numpy as np from scipy import stats def detect_drift(train_data, online_data, threshold0.2): # 计算PSI def psi(expected, actual, buckets10): breakpoints np.percentile(expected, np.linspace(0, 100, buckets 1)) expected_counts np.histogram(expected, breakpoints)[0] / len(expected) actual_counts np.histogram(actual, 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) return np.sum((actual_counts - expected_counts) * np.log(actual_counts / expected_counts)) psi_value psi(train_data, online_data) return psi_value threshold, psi_value6.2 反馈闭环bad case 是最好的老师监控发现问题之后怎么闭环我的做法是建立一个bad case 收集管道。线上预测置信度低、或者用户反馈错误的样本自动收集起来定期人工复核确认是数据问题还是模型问题。如果是数据问题修数据管道。如果是模型问题把这些样本加入训练集重新训练。这个循环跑起来之后模型会持续进化。反馈闭环的关键是速度。从发现问题到模型更新上线周期越短越好。我见过团队一个月才更新一次模型结果线上问题拖了一个月才解决。理想情况下这个周期应该是一周以内甚至更短。实现快速迭代的前提是自动化。数据收集自动化、训练自动化、评估自动化、部署自动化。每一步都自动化之后人工只需要做决策——要不要上线这个新版本。6.3 再训练策略什么时候该更新模型再训练不是越频繁越好。频繁再训练成本高而且可能引入不稳定性。我的策略是触发式再训练满足以下任一条件就触发数据漂移超过阈值PSI 0.2模型指标下降超过阈值比如AUC下降5%积累了一定量的新标注数据比如1万条固定周期比如每月一次作为兜底再训练之后不能直接全量上线要先做影子模式——新模型和旧模型同时跑新模型只记录不生效对比两者的表现。确认新模型更好之后再逐步切流量比如先切10%观察一天没问题再切50%最后全量。这个流程看起来繁琐但能避免“新模型上线导致业务指标暴跌”这种事故。我亲身经历过一次新模型离线指标很好上线后转化率掉了15%查了半天发现是特征计算的一个边界条件没处理好。从那以后我再也不敢跳过影子模式。7. 常见问题与排查技巧实录7.1 训练时loss不下降怎么排查这是最常见的问题。我的排查顺序是检查数据标签对不对特征有没有归一化有没有NaN检查模型输出层激活函数对不对损失函数选对了吗检查优化器学习率是不是太大或太小试几个数量级。检查梯度有没有梯度消失或爆炸打印梯度范数看看。学习率是最常见的元凶。我习惯先用一个较大的学习率跑几步看loss有没有下降然后逐步调小。如果loss完全不降大概率是数据或模型结构的问题。7.2 线上效果比离线差很多怎么回事这个问题我遇到过好几次原因通常是三类训练/推理不一致特征计算逻辑不一致或者预处理步骤漏了。数据泄漏训练时用了未来信息线上拿不到。分布偏移线上数据分布和训练数据不一样。排查方法是在线上采样一批数据用训练时的代码重新算一遍特征对比线上实际用的特征。如果对不上就是一致性问题。如果对得上但效果还是差那就是分布偏移需要重新采样训练数据。7.3 推理服务延迟高怎么优化延迟优化按优先级排序优化手段预期收益实施难度批处理高中模型量化中高低算子融合中低缓存高命中时低模型蒸馏高高硬件升级高低但贵我一般先上批处理和量化这两个投入产出比最高。缓存适合重复请求多的场景。模型蒸馏和硬件升级是最后的手段。7.4 模型版本回滚怎么做到快速安全回滚的关键是版本管理规范和部署自动化。每个模型版本都有完整的环境和依赖回滚时直接切到旧版本的镜像不需要重新构建。我习惯保留最近5个版本的热备更早的版本归档到冷存储。回滚操作应该是一键式的最好能在1分钟内完成。如果回滚需要半小时那说明部署流程有问题得先修流程。提示回滚之后一定要复盘搞清楚为什么新版本有问题。否则下次还会踩同样的坑。8. 我踩过的几个大坑希望你别再踩第一个坑是过早优化。项目刚开始就想着上分布式、上GPU集群结果数据量根本撑不起来白白浪费几周时间搭环境。后来我学乖了先用单机跑通全流程等真的遇到瓶颈再优化。单机Pandas能处理千万级数据单卡GPU能训练大多数中小模型这些足够验证想法了。第二个坑是忽略数据质量。有次模型效果怎么调都上不去折腾了两周最后发现是数据里有一批标签标反了。清理之后效果直接达标。从那以后我每次新数据集进来第一件事就是抽样人工检查确认数据质量没问题再往下走。第三个坑是监控缺失。早期上线没有监控模型效果下降了都不知道等业务方反馈才发现。后来补上了监控但只监控了系统指标没监控模型指标。结果数据漂移了三个月才发现期间业务损失不小。现在我的标准是没有监控的模型不允许上线。第四个坑是文档缺失。有次排查一个线上问题发现是半年前的一个特征处理逻辑有问题但当时写代码的人已经离职了文档也没写只能对着代码一行行猜。从那以后我要求每个特征、每个模型、每个服务都必须有文档写清楚输入输出、依赖关系、已知问题。这些坑说到底都是工程规范的问题。AI工程和传统软件工程没有本质区别规范、测试、监控、文档一个都不能少。模型再先进工程跟不上照样落不了地。最后分享一个我常用的检查清单每次上线前过一遍数据版本固定了吗训练/推理特征一致吗模型版本完整吗权重结构预处理环境监控埋点了吗降级逻辑有了吗回滚方案验证了吗文档写了吗这七个问题都能答上来上线基本不会出大问题。答不上来那就再等等。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑