☰
从零手搓AI工程化全流程:数据、训练、推理与监控实战
2026/9/30 4:39:11 网站建设 项目流程

1. 为什么我要从零手搓一套 AI 工程化流程

第一次看到ai-engineering-from-scratch这个项目名的时候,我正坐在工位上啃一个已经迭代了三个版本的推荐系统。那套系统用了大量现成的框架和封装好的推理服务,跑是能跑,但每次线上出问题,排查链路长得让人崩溃——数据预处理在 A 服务,特征拼接在 B 服务,模型推理在 C 服务,后处理又在 D 服务,日志散落在四个地方,一个简单的"为什么这条样本打分异常"能查一整个下午。

那一刻我突然意识到一个问题:我们这行太多人只会"调包",不会"造轮子"。会用transformers加载模型,会用FastAPI起个接口,会用Docker打个镜像,但真让你从一张白纸开始,把数据加载、模型训练、推理服务、监控告警这一整条链路自己搭一遍,很多人是懵的。ai-engineering-from-scratch这个项目标题戳中的正是这个痛点——从零开始,把 AI 工程化的每一个环节都亲手实现一遍。

这篇文章不是教你调某个库的 API,而是把我自己从零搭建一套 AI 工程化流程的完整思路、踩过的坑、以及那些文档里不会写的经验,原原本本分享出来。适合谁看?如果你已经会写 Python、懂一点机器学习基础,但每次做项目都停留在"跑通 notebook"的阶段,不知道怎么把它变成一个能上线、能维护、能排查问题的工程系统,那这篇内容就是写给你的。如果你已经是资深工程师,也可以看看我在选型和架构上的取舍逻辑,说不定能给你一些不同角度的参考。

我打算按这么几个部分来讲:先讲整体设计思路和方案选型背后的考量,再拆核心细节和实操要点,然后是完整的实操过程和关键环节实现,最后是我在实际操作中遇到的典型问题和排查技巧。全程都是我自己趟出来的路,能抄作业的地方我会把参数和命令都给全。

2. 整体设计与思路拆解

2.1 为什么选择"从零实现"而不是"拼装现成组件"

很多人会问,现在生态这么成熟,LangChain、LlamaIndex、vLLM、Triton这些工具拿来就用,为什么还要从零写?我的答案很直接:拼装能让你快速出活,但从零实现能让你真正理解系统。

我举个具体的例子。之前团队里有个小伙子,用现成的推理框架部署了一个文本分类模型,QPS 上不去,他调了两天参数都没解决。后来我让他把整个推理链路拆开看,才发现瓶颈根本不在模型本身,而在 tokenizer 那一层——他用的分词器是 Python 实现的,单线程处理,每次请求都要重新加载词表。这个问题在封装好的框架里被藏得很深,只有你自己实现过一遍,才知道每一层的时间开销在哪里。

从零实现还有几个实际好处。第一,依赖可控。现成框架版本一升级,API 一变动,你的代码可能就跑不起来了,而从零写的核心逻辑,依赖的只有最基础的库,稳定性高得多。第二,可定制性强。业务需求千奇百怪,框架的抽象层有时候反而是束缚,自己写的代码想怎么改就怎么改。第三,排查问题快。整条链路都是你写的,出了问题你知道该去哪里看,不用去翻别人的源码。

当然,我不是说所有场景都要从零写。快速验证想法的时候,用现成工具没问题。但如果你想真正掌握 AI 工程化这门手艺,从零实现一遍是绕不过去的坎。

2.2 整体架构的分层设计

我把整套系统分成了四层,这个分层是我踩了很多坑之后定下来的,每一层的职责边界都很清晰。

第一层是数据层。负责原始数据的读取、清洗、切分、缓存。这一层的核心目标是"把脏活累活都干完",让上层拿到的是干净、规整的数据。我见过太多项目把数据清洗逻辑散落在训练脚本里,结果训练和推理时的预处理不一致,线上效果直接崩掉。数据层单独抽出来,训练和推理共用同一套预处理逻辑,这个问题就从根本上避免了。

第二层是模型层。负责模型的定义、训练、评估、保存。这一层的关键是"可复现"——同样的数据、同样的配置,跑出来的结果必须一致。我强制要求所有随机种子都要固定,所有超参数都要写进配置文件,不允许硬编码在代码里。

第三层是服务层。负责把训练好的模型包装成可调用的服务。这一层要考虑的是并发、批处理、超时、降级这些工程问题。我选择用异步框架来做,因为 AI 推理本质上是 IO 密集和计算密集混合的场景,异步能显著提升吞吐。

第四层是监控层。负责记录每一次请求的输入输出、耗时、异常。这一层最容易被忽略,但恰恰是线上系统能不能维护好的关键。我的原则是:任何一次推理调用,都必须留下可追溯的记录。

这四层之间通过明确的接口通信,层与层之间不直接依赖具体实现。比如数据层对外只暴露"给我一批处理好的样本"这个接口,至于底层是从 CSV 读还是从数据库读,上层不关心。这种设计让每一层都可以独立替换和测试。

2.3 技术选型的取舍逻辑

选型这块我纠结了很久,最后定下来的方案是:核心逻辑用纯 Python 实现,只在真正需要性能的地方引入第三方库。

数据处理我用pandas和numpy,这两个是事实标准,没必要自己造。但数据切分和缓存的逻辑我自己写,因为这部分逻辑和业务强相关,用现成的反而别扭。模型训练我用PyTorch,这个也没得选,生态太成熟了。但训练循环我自己写,不用Trainer那套封装,因为我想完全掌控每一个 step 里发生了什么。

服务层我用FastAPI加uvicorn,异步支持好,性能也够。批处理逻辑我自己实现,因为现成的批处理方案要么太重,要么不满足我的需求。监控层我用logging加SQLite,简单直接,不引入额外的中间件依赖。有人会问为什么不用Prometheus加Grafana,我的考虑是:在项目早期,简单可靠比功能强大更重要。SQLite 单文件,部署零成本,查询用 SQL 就行,等业务量真上来了再换也不迟。

这里有个选型原则我想强调一下:不要为了用某个工具而用某个工具。我见过太多项目上来就堆一整套可观测性栈,结果团队里没人会维护,出了问题反而更麻烦。工具是为人服务的,选你能 hold 住的。

3. 核心细节解析与实操要点

3.1 数据层的预处理一致性保障

数据层最容易出问题的地方,就是训练和推理的预处理不一致。我举个真实的例子:训练的时候文本做了小写转换,推理的时候忘了做,结果同一个句子,训练时模型见过的是小写版本,推理时喂进去的是原始大小写混合版本,效果直接掉一大截。这种问题非常隐蔽,因为代码不报错,只是效果变差。

我的解决方案是把预处理逻辑封装成一个独立的类,训练和推理都调用同一个实例。这个类对外只暴露两个方法:fit和transform。fit负责从训练数据里学习必要的统计量,比如词表、归一化参数;transform负责把原始数据转成模型能吃的格式。训练时先fit再transform,推理时直接加载保存好的状态再transform。

class TextPreprocessor: def __init__(self, lowercase=True, max_len=128): self.lowercase = lowercase self.max_len = max_len self.vocab = None def fit(self, texts): # 从训练数据构建词表 processed = [self._normalize(t) for t in texts] self.vocab = self._build_vocab(processed) def transform(self, texts): if self.vocab is None: raise RuntimeError("必须先调用 fit 或加载已保存的状态") result = [] for t in texts: normalized = self._normalize(t) ids = self._to_ids(normalized) result.append(self._pad_or_truncate(ids)) return result def _normalize(self, text): if self.lowercase: text = text.lower() return text.strip()

这里有个细节要注意:fit学到的状态必须能序列化保存。我用pickle存,简单可靠。保存的时候连同版本号一起存,加载的时候校验版本,避免新旧版本不兼容导致静默错误。

提示:预处理类的状态一定要版本化。我吃过亏,改了预处理逻辑但忘了更新版本号,结果线上加载了旧状态,排查了半天才发现。

3.2 模型训练的可复现性设计

可复现性是 AI 工程化的生命线。一个不能复现的实验,等于没做。我要求所有训练必须满足三个条件:随机种子固定、配置外置、日志完整。

随机种子这块,光设torch.manual_seed是不够的。numpy、random、甚至 CUDA 的随机数都要设。我写了个工具函数统一处理:

import random import numpy as np import torch def set_seed(seed=42): random.seed(seed) np.random.seed(seed) torch.manual_seed(seed) torch.cuda.manual_seed_all(seed) torch.backends.cudnn.deterministic = True torch.backends.cudnn.benchmark = False

注意最后两行,deterministic=True会让 cuDNN 只用确定性的算法,benchmark=False关掉自动调优。这两个设置会牺牲一点性能,但换来的是完全可复现。如果你的场景对性能极度敏感,可以只在调试阶段开,生产环境关掉。

配置外置我用 YAML 文件,所有超参数、路径、模型结构都写在里面。代码里不允许出现魔法数字。这样做的好处是,每次实验的配置都能存档,想复现哪个实验,把对应的配置文件拿出来就行。

# config/train_base.yaml seed: 42 data: train_path: data/train.csv val_path: data/val.csv max_len: 128 model: vocab_size: 30000 embed_dim: 256 num_heads: 8 num_layers: 4 train: batch_size: 64 lr: 0.0001 epochs: 10 warmup_steps: 1000

日志这块,我要求每个 epoch 结束都要记录:训练损失、验证损失、学习率、耗时、显存占用。这些数据后面分析训练曲线、排查过拟合都用得上。我用logging写到文件,同时打印到控制台。

3.3 服务层的批处理与并发控制

服务层最核心的两个问题:批处理和并发控制。这两个问题处理不好,要么吞吐上不去,要么延迟爆炸。

批处理的思路是:请求进来先不急着推理,放进一个队列,攒够一定数量或者等一小段时间,凑成一批一起推理。这样能充分利用 GPU 的并行能力。但这里有个权衡:攒批会增加延迟,因为请求要等。我的做法是设置两个阈值——最大批大小和最大等待时间,哪个先到就触发推理。

import asyncio from collections import deque class BatchScheduler: def __init__(self, max_batch_size=32, max_wait_ms=50): self.max_batch_size = max_batch_size self.max_wait = max_wait_ms / 1000 self.queue = deque() self.lock = asyncio.Lock() async def submit(self, item): future = asyncio.get_event_loop().create_future() async with self.lock: self.queue.append((item, future)) if len(self.queue) >= self.max_batch_size: await self._flush() # 等待结果 return await future async def _flush(self): if not self.queue: return batch = list(self.queue) self.queue.clear() items = [b[0] for b in batch] futures = [b[1] for b in batch] try: results = await self._infer_batch(items) for fut, res in zip(futures, results): fut.set_result(res) except Exception as e: for fut in futures: fut.set_exception(e)

并发控制这块,我用asyncio.Semaphore限制同时进行的推理数量。因为 GPU 显存有限,同时跑太多批会 OOM。信号量的值根据显存大小和单批显存占用算出来。比如显存 16G,单批占用 2G,留点余量,信号量设 6 比较稳妥。

注意:批处理的最大等待时间不要设太长。我一开始设了 200ms,结果 P99 延迟直接飙到 300ms,用户体验很差。后来改成 50ms,吞吐只降了不到 10%,但延迟降了一半多。

3.4 监控层的埋点设计

监控层的核心是埋点。埋点埋得好,线上问题一目了然;埋得不好,等于没埋。我的埋点原则是:记录一切可能出问题的环节。

具体来说,每次请求我记录这些字段:请求 ID、时间戳、输入长度、输出长度、推理耗时、排队耗时、批大小、是否命中缓存、异常信息。这些字段存进 SQLite,用请求 ID 关联。

import sqlite3 import time import uuid class RequestLogger: def __init__(self, db_path="logs/requests.db"): self.conn = sqlite3.connect(db_path, check_same_thread=False) self._init_table() def _init_table(self): self.conn.execute(""" CREATE TABLE IF NOT EXISTS requests ( request_id TEXT PRIMARY KEY, timestamp REAL, input_len INTEGER, output_len INTEGER, infer_ms REAL, queue_ms REAL, batch_size INTEGER, cache_hit INTEGER, error TEXT ) """) self.conn.commit() def log(self, **kwargs): kwargs.setdefault("request_id", str(uuid.uuid4())) kwargs.setdefault("timestamp", time.time()) cols = ", ".join(kwargs.keys()) placeholders = ", ".join(["?"] * len(kwargs)) self.conn.execute( f"INSERT INTO requests ({cols}) VALUES ({placeholders})", list(kwargs.values()) ) self.conn.commit()

有了这些数据,排查问题就方便多了。比如发现某段时间延迟高,直接查queue_ms和infer_ms,就能判断是排队排太久还是推理本身慢。再比如发现效果变差,查input_len的分布,看看是不是输入长度分布变了。

这里有个经验:埋点不要怕多,怕的是不够。多记几个字段,存储成本增加有限,但排查问题时能省大量时间。等系统稳定了,再考虑精简。

4. 实操过程与核心环节实现

4.1 环境准备与依赖管理

环境这块我强烈建议用虚拟环境,不要污染全局 Python。我用venv,简单够用。依赖管理用requirements.txt,但我会把直接依赖和间接依赖分开,直接依赖手动维护,间接依赖用pip freeze生成。

python -m venv venv source venv/bin/activate # Windows 用 venv\Scripts\activate pip install torch numpy pandas fastapi uvicorn pyyaml pip freeze > requirements-lock.txt

requirements.txt里我只写直接依赖,不锁版本,方便升级。requirements-lock.txt锁死所有版本,用于生产部署。这样开发时灵活,部署时稳定。

提示:torch的安装要特别注意 CUDA 版本匹配。我踩过坑,装了个 CPU 版本的 torch,训练慢得让人怀疑人生。装之前先nvidia-smi看 CUDA 版本,然后去官网找对应的安装命令。

4.2 数据管道的搭建

数据管道我分成了三个步骤:加载、清洗、切分。加载用pandas读 CSV,清洗包括去重、去空、长度过滤,切分按 8:1:1 分训练、验证、测试。

import pandas as pd from sklearn.model_selection import train_test_split def build_dataset(raw_path, output_dir, seed=42): df = pd.read_csv(raw_path) # 去重 df = df.drop_duplicates(subset=["text"]) # 去空 df = df.dropna(subset=["text", "label"]) # 长度过滤 df["text_len"] = df["text"].str.len() df = df[(df["text_len"] >= 2) & (df["text_len"] <= 512)] # 切分 train, temp = train_test_split(df, test_size=0.2, random_state=seed, stratify=df["label"]) val, test = train_test_split(temp, test_size=0.5, random_state=seed, stratify=temp["label"]) train.to_csv(f"{output_dir}/train.csv", index=False) val.to_csv(f"{output_dir}/val.csv", index=False) test.to_csv(f"{output_dir}/test.csv", index=False) print(f"训练集 {len(train)} 条,验证集 {len(val)} 条,测试集 {len(test)} 条")

这里有个细节:切分的时候用了stratify,保证各个集合里标签分布一致。这个很重要,尤其是类别不平衡的数据,不 stratify 的话可能某个类别在验证集里一条都没有,评估结果完全不可信。

数据清洗的阈值怎么定?长度下限 2 是拍脑袋定的,因为太短的文本没信息量。上限 512 是根据模型的最大长度定的,超过的截断,但截断会丢信息,所以干脆过滤掉。这些阈值没有标准答案,要根据你的数据分布来定。我的建议是先画个长度分布直方图,看看 95 分位数在哪里,上限就定在那里附近。

4.3 模型训练循环的实现

训练循环我自己写,核心就几个部分:前向、算损失、反向、更新、评估。但每个部分都有讲究。

def train_one_epoch(model, dataloader, optimizer, scheduler, device, epoch): model.train() total_loss = 0 start = time.time() for step, batch in enumerate(dataloader): input_ids = batch["input_ids"].to(device) labels = batch["labels"].to(device) optimizer.zero_grad() logits = model(input_ids) loss = F.cross_entropy(logits, labels) loss.backward() # 梯度裁剪,防止梯度爆炸 torch.nn.utils.clip_grad_norm_(model.parameters(), max_norm=1.0) optimizer.step() scheduler.step() total_loss += loss.item() if step % 100 == 0: lr = scheduler.get_last_lr()[0] print(f"Epoch {epoch} Step {step} Loss {loss.item():.4f} LR {lr:.6f}") avg_loss = total_loss / len(dataloader) elapsed = time.time() - start print(f"Epoch {epoch} 完成,平均损失 {avg_loss:.4f},耗时 {elapsed:.1f}s") return avg_loss

梯度裁剪这步很关键。我遇到过训练到一半 loss 突然变成 NaN 的情况,就是因为某一步梯度太大。加上clip_grad_norm_之后,这个问题基本没再出现过。max_norm设 1.0 是个经验值,大部分场景够用,如果你的模型特别深,可以适当调大。

学习率调度我用的是 warmup 加 cosine 衰减。warmup 让学习率从 0 慢慢升到设定值,避免一开始步子太大把模型带偏。cosine 衰减让学习率平滑下降,比阶梯式下降效果通常更好。

from torch.optim.lr_scheduler import LambdaLR import math def get_cosine_schedule(optimizer, warmup_steps, total_steps): def lr_lambda(step): if step < warmup_steps: return step / max(1, warmup_steps) progress = (step - warmup_steps) / max(1, total_steps - warmup_steps) return max(0.0, 0.5 * (1 + math.cos(math.pi * progress))) return LambdaLR(optimizer, lr_lambda)

评估的时候一定要model.eval()加torch.no_grad(),前者关掉 dropout 和 batchnorm 的训练行为,后者省显存。这两个忘了任何一个,评估结果都不准。

4.4 推理服务的完整实现

推理服务我用FastAPI搭,核心是一个异步的推理端点。请求进来先过预处理,然后提交给批调度器,等结果返回后再做后处理。

from fastapi import FastAPI, HTTPException from pydantic import BaseModel import asyncio app = FastAPI() scheduler = BatchScheduler(max_batch_size=32, max_wait_ms=50) preprocessor = None model = None class InferRequest(BaseModel): text: str class InferResponse(BaseModel): label: int score: float request_id: str @app.post("/predict", response_model=InferResponse) async def predict(req: InferRequest): if not req.text.strip(): raise HTTPException(status_code=400, detail="文本不能为空") request_id = str(uuid.uuid4()) start = time.time() try: # 预处理 input_ids = preprocessor.transform([req.text])[0] # 提交批调度 result = await scheduler.submit(input_ids) # 后处理 label = int(result.argmax()) score = float(result.softmax(dim=-1).max()) elapsed = (time.time() - start) * 1000 logger.log( request_id=request_id, input_len=len(req.text), output_len=1, infer_ms=elapsed, batch_size=scheduler.last_batch_size, cache_hit=0, error=None ) return InferResponse(label=label, score=score, request_id=request_id) except Exception as e: logger.log(request_id=request_id, error=str(e)) raise HTTPException(status_code=500, detail="推理失败")

启动命令:

uvicorn app:app --host 0.0.0.0 --port 8000 --workers 1

注意workers设 1。因为批调度器是进程内的,多 worker 会导致每个 worker 各自维护一个队列,批处理效果打折。如果要多进程,得把批调度器抽出来做成独立的服务,用消息队列通信。这个复杂度在项目早期没必要,单 worker 加异步已经能扛不少量了。

4.5 监控数据的查询与分析

监控数据存进 SQLite 之后,用 SQL 查询就行。我常用的几个查询:

-- 最近一小时的 P99 延迟 SELECT infer_ms FROM requests WHERE timestamp > strftime('%s', 'now', '-1 hour') ORDER BY infer_ms DESC LIMIT 1 OFFSET (SELECT COUNT(*) * 0.01 FROM requests WHERE timestamp > strftime('%s', 'now', '-1 hour')); -- 各批大小的请求分布 SELECT batch_size, COUNT(*) as cnt FROM requests GROUP BY batch_size ORDER BY cnt DESC; -- 异常请求 SELECT request_id, error FROM requests WHERE error IS NOT NULL;

这些查询能回答大部分线上问题。比如 P99 延迟高,看是不是某个批大小特别慢;异常多,看错误信息是什么。我建议把这些查询封装成脚本,出问题的时候一键跑,省得临时写 SQL。

5. 常见问题与排查技巧实录

5.1 训练不收敛的排查思路

训练不收敛是最常见的问题,表现是 loss 不降或者震荡。我的排查顺序是这样的:

第一步,检查数据。把一批数据打印出来,看看输入和标签对不对。我遇到过标签错位的情况,输入是第 i 条,标签是第 i+1 条,这种问题不看数据根本发现不了。

第二步,检查学习率。学习率太大 loss 会震荡甚至爆炸,太小 loss 降得极慢。我的经验是从 1e-4 开始试,如果 loss 不动就调大到 1e-3,如果震荡就调小到 1e-5。

第三步,检查模型。用一个极小的数据集(比如 10 条)过拟合一下,如果连 10 条都拟合不了,说明模型结构有问题。这个技巧非常有用,能快速定位是数据问题还是模型问题。

第四步,检查梯度。打印每一层的梯度范数,看看有没有梯度消失或爆炸。如果某一层梯度一直是 0,可能是激活函数选错了,或者初始化有问题。

5.2 推理延迟高的定位方法

推理延迟高,先看监控数据里的queue_ms和infer_ms。如果queue_ms高,说明请求排队严重,要么是并发太高,要么是批处理等待时间太长。如果infer_ms高,说明推理本身慢,可能是批太大或者模型太重。

我遇到过一次延迟突然翻倍的情况,查监控发现infer_ms没变,但queue_ms涨了 10 倍。进一步查发现是某个客户端疯狂发请求,把队列撑爆了。解决办法是加限流,每个客户端每秒最多 100 个请求。

还有一种情况是显存碎片化导致推理变慢。这个比较隐蔽,表现是跑一段时间后延迟逐渐升高。解决办法是定期重启服务,或者用torch.cuda.empty_cache()手动清理。不过empty_cache会带来短暂卡顿,要挑低峰期做。

5.3 常见问题速查表

问题现象可能原因排查方法解决方案
训练 loss 不降学习率太小、数据有问题、模型结构错小数据集过拟合测试调学习率、检查数据、改模型
训练 loss 震荡学习率太大、批太小打印每步 loss调小学习率、增大批
推理延迟高排队久、批太大、显存碎片看 queue_ms 和 infer_ms限流、调批大小、定期重启
线上效果差预处理不一致、数据分布漂移对比训练和推理预处理统一预处理、重新训练
显存 OOM批太大、模型太大、内存泄漏看显存占用曲线调小批、梯度累积、查泄漏
服务启动失败端口占用、依赖缺失、配置错看启动日志换端口、装依赖、改配置

5.4 几个我踩过的坑

坑一:pickle 加载预处理状态时路径不对。我习惯用相对路径,结果服务从不同目录启动时找不到文件。后来改成基于文件位置计算绝对路径,问题解决。

坑二:异步代码里调了同步阻塞函数。我在推理函数里直接调了model(input),这是同步的,会阻塞事件循环,导致并发上不去。后来改成用run_in_executor包一层,或者用支持异步的推理库。

坑三:日志写 SQLite 并发冲突。多个协程同时写 SQLite 会报database is locked。解决办法是加锁,或者用 WAL 模式。我用的是加锁,简单可靠。

import threading self.write_lock = threading.Lock() def log(self, **kwargs): with self.write_lock: # 写数据库 ...

坑四:模型保存时忘了保存配置。只存了权重,加载的时候不知道模型结构是什么,只能靠记忆。后来我改成权重和配置一起存,加载时先读配置再建模型。

提示:模型保存建议用torch.save({"state_dict": model.state_dict(), "config": config})这种格式,把配置一起存进去,加载时就不会抓瞎。

6. 后续可以扩展的方向

这套系统跑通之后,我陆续加了一些扩展。第一个是缓存,对相同的输入直接返回缓存结果,省掉推理开销。缓存用 LRU 策略,键是输入的哈希值。这个对重复请求多的场景效果显著,我实测命中率能到 30% 左右。

第二个是A/B 测试,支持同时部署多个模型版本,按流量比例分流。这个对模型迭代很有用,新模型先小流量验证,没问题再全量。

第三个是自动扩缩容,根据队列长度动态调整 worker 数量。这个复杂度比较高,我还在摸索阶段,等成熟了再单独写一篇。

这套东西我从零搭起来花了大概两周,中间踩的坑比写代码的时间还多。但搭完之后,再遇到 AI 工程化的问题,我心里就有底了,因为每一层我都亲手实现过,知道问题可能出在哪里。这种掌控感,是调包调不出来的。如果你也想真正掌握 AI 工程化,我建议你找个周末,从零搭一遍,哪怕功能简单点,收获也会很大。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询