1. 这不是调包,是亲手搭起AI工程的钢筋水泥
“AI Engineering from Scratch”——看到这个标题,我第一反应不是兴奋,而是下意识摸了摸键盘边沿那层被磨得发亮的漆。过去三年,我带过17个团队落地AI项目,从智能客服到工业质检,从医疗影像辅助标注到供应链需求预测。几乎每个新来的工程师开口第一句都是:“老师,用哪个框架?Hugging Face Model Hub里挑一个预训练模型行不行?”——行,当然行,但行完之后呢?模型在测试集上AUC 0.92,上线后第二天监控告警疯狂闪烁,延迟从200ms飙到3.8秒,错误率翻了4倍。没人知道为什么。因为没人真正“从头”看过它怎么呼吸、怎么吃数据、怎么排泄日志、怎么在内存里喘息。
AI Engineering from Scratch,不是教你怎么从零写Transformer,也不是让你手搓CUDA核函数。它是一套可交付、可运维、可演进的AI系统建造方法论:从你敲下第一个mkdir ai-engineering开始,到最终交付一个能扛住每秒500次并发请求、自动降级、带全链路追踪、模型版本可回滚、数据漂移能预警的生产级服务为止。它解决的不是“能不能跑”,而是“敢不敢上线”“出了问题找谁”“下周业务要加个新字段怎么改不崩”。关键词ai-engineering指的是把AI当作一个需要持续集成、持续部署、可观测、有SLO保障的软件子系统来对待;而from-scratch强调的是一种“拆解—重建—验证”的思维惯性——不依赖黑盒封装,不跳过任何中间态,每一个模块都经得起拷问:它的输入边界在哪?失败时如何兜底?资源消耗是否可量化?日志是否能定位到具体batch?
适合谁?不是纯算法研究员,也不是只会写API的后端。而是那些卡在“模型训好了,但不知道下一步该交给谁”的ML工程师;是接到需求后第一反应不是查文档而是画架构图的Tech Lead;是运维同事半夜打电话问“你们那个模型服务占了80%内存,到底在干啥”的背锅侠;更是想真正理解AI系统全貌、摆脱“调参炼丹师”身份的技术负责人。它不承诺让你三天写出LLM,但它保证:当你下次再看到“OOM Killed”报错时,你能立刻判断是模型加载策略问题,还是批处理队列堆积,而不是重启服务器碰运气。
我试过用现成MLOps平台快速交付——确实快,两周上线。但第三个月,业务方要求把输入字段从JSON改成Protobuf序列化,平台不支持;第四个月,模型需要接入内部认证网关,插件机制僵硬;第六个月,监控发现某类样本推理耗时突增300%,平台只给个平均P95,查不出是哪一层算子拖慢。最后我们花了三周重写整个服务骨架——这次,我们从Scratch开始。
2. 整体设计:拒绝“模型即服务”,构建四层可拆卸架构
AI Engineering from Scratch 的核心,是把一个看似原子的“AI服务”拆解为四个物理隔离、职责清晰、可独立演进的层次。这不是理论分层,而是我们踩坑后用血泪画出的边界线。每一层都必须能单独测试、单独部署、单独扩缩容,且层间通信必须通过明确定义的契约(Schema + Protocol)。
2.1 第一层:数据摄取与预处理管道(Data Ingestion & Preprocessing Pipeline)
这是整个系统的“消化系统”。很多人以为AI工程从模型开始,其实它从数据如何进入系统开始。我们不用Airflow或Prefect这类通用编排工具做第一层,原因很现实:它们太重,且对实时性、schema变更、错误恢复的支持过于抽象。我们选择轻量级、可嵌入的方案——用Rust写的>rule_id: "high-risk-override" condition: "model_score > 0.85 && user_tenure_months < 3" action: "block_transaction"Orchestrator启动时加载所有规则,每次请求先调用模型,再根据结果+原始业务上下文(来自HTTP Header或额外字段)执行规则匹配。规则变更无需重启服务,Consul Watch自动热更新。
2.4 第四层:可观测性与反馈闭环(Observability & Feedback Loop)
AI系统最怕“黑盒静默失效”。我们投入30%开发时间建这一层,因为它决定系统能否长期存活。
- 三维监控体系:
- 基础设施层:CPU/Mem/Disk IO(Prometheus + Grafana),重点看
model-runner进程RSS内存增长曲线——缓慢爬升意味着内存泄漏。 - 服务层:gRPC指标(成功率、P95延迟、QPS)、Sidecar指标(重试率、TLS握手失败)、Orchestrator指标(规则匹配率、Saga补偿次数)。
- AI层:模型输入分布(Histogram of
model_score)、输出置信度分布、特征漂移检测(Evidently.ai计算PSI,阈值>0.15告警)、概念漂移(在线计算KL散度,突增触发人工审核)。
- 基础设施层:CPU/Mem/Disk IO(Prometheus + Grafana),重点看
- 自动反馈闭环:所有线上预测结果(含用户最终操作:接受/拒绝/申诉)存入Delta Lake。每小时触发Spark Job,计算模型在各人群上的F1-score衰减率。若某人群F1下降>15%,自动生成Jira Ticket并附上典型bad case样本,分配给对应算法工程师。
- 人工干预通道:Orchestrator提供REST API
/override,运营人员可手动覆盖模型决策(如“此订单强制通过”),所有覆盖记录存入审计日志,并触发模型重训数据采样——覆盖行为本身成为新训练数据的强信号。
这套四层架构不是一步到位。我们第一版只有两层(模型服务+简单API),第二版加了Orchestrator,第三版才补全可观测性。每次迭代都源于一次线上事故:第一次OOM让我们加进程隔离,第一次规则变更引发资损让我们拆出Orchestrator,第一次模型静默退化让我们建反馈闭环。From Scratch的本质,是让每个模块的诞生都有明确的“痛感驱动”。
3. 核心细节:从代码到部署的12个关键实操点
光有架构不够,细节决定生死。以下是我们在真实项目中反复验证、写进团队SOP的12个核心实操点,每个都附带“为什么”和“怎么做”。
3.1 数据校验:用Avro Schema而非JSON Schema
为什么?JSON Schema无法描述二进制数据(如图像base64)、无法定义浮点数精度、不支持向后兼容演进。Avro Schema强制要求字段有默认值或标记为null,且版本升级时遵循严格兼容规则(添加可选字段OK,删除字段NG)。
怎么做?
- 定义Schema(
user_event.avsc):{ "type": "record", "name": "UserEvent", "fields": [ {"name": "event_id", "type": "string"}, {"name": "timestamp", "type": "long"}, {"name": "image_data", "type": ["null", "bytes"], "default": null}, {"name": "score", "type": "double", "default": 0.0} ] } - 在
>val schema = new Schema.Parser().parse(new File("user_event.avsc")) val decoder = DecoderFactory.get().binaryDecoder(inputStream, null) val reader = new GenericDatumReader[GenericRecord](schema) val record = reader.read(null, decoder) // 自动校验类型与必填
注意:Avro Schema必须托管在Confluent Schema Registry或自建Avro Registry,禁止本地文件引用。否则不同服务加载不同版本Schema,数据就乱了。
3.2 模型加载:避免PyTorch的torch.load()直接反序列化
为什么?torch.load()会执行任意Python代码,存在远程代码执行风险;且加载大模型(>2GB)时阻塞主线程,服务启动慢。
怎么做?
- 使用
torch.jit.script导出模型:model = MyModel() model.eval() traced_model = torch.jit.script(model) traced_model.save("model.pt") # 生成纯二进制,无Python字节码 model-runner用torch.jit.load()加载:import torch model = torch.jit.load("model.pt") # 安全、快速、线程安全 model.to("cuda:0") # 显存分配在此刻完成,非首次推理时
3.3 内存管理:为每个model-runner设置cgroup内存限制
为什么?PyTorch CUDA缓存(torch.cuda.memory_cached())不会随模型销毁自动释放,多个模型实例共用GPU时,内存碎片化严重,最终OOM。
怎么做?
- 启动脚本中设置cgroup:
# 创建cgroup sudo cgcreate -g memory:/model-fraud-v3.2 # 设置内存上限2GB echo 2147483648 | sudo tee /sys/fs/cgroup/memory/model-fraud-v3.2/memory.limit_in_bytes # 启动进程 sudo cgexec -g memory:model-fraud-v3.2 ./model-runner --model-path s3://... - 在
model-runner中,定期清理缓存:import torch def cleanup_gpu_memory(): if torch.cuda.is_available(): torch.cuda.empty_cache() # 清理缓存 torch.cuda.synchronize() # 确保同步 # 每100次推理后调用
3.4 推理超时:gRPC层面设timeout,而非模型内设
为什么?模型内部设timeout(如torch.set_num_threads(1))无法中断CUDA核函数;而gRPC timeout可在网络层强制断开连接,触发Sidecar重试。
怎么做?
model-runnergRPC Server配置:server = grpc.server( futures.ThreadPoolExecutor(max_workers=10), options=[ ('grpc.max_send_message_length', 100 * 1024 * 1024), # 100MB ('grpc.max_receive_message_length', 100 * 1024 * 1024), ] ) # 关键:设置每个RPC的deadline server.add_insecure_port('[::]:50051')- Client调用时显式设timeout:
try: response = stub.Predict(request, timeout=5.0) # 5秒硬超时 except grpc.RpcError as e: if e.code() == grpc.StatusCode.DEADLINE_EXCEEDED: # 触发降级逻辑 return fallback_result()
3.5 特征漂移检测:用PSI而非KS检验
为什么?KS检验对样本量敏感,小样本(<1000)易误报;PSI(Population Stability Index)对分布变化更鲁棒,且可解释性强(PSI>0.1警告,>0.25严重)。
怎么做?
- 计算PSI公式:
PSI = Σ (Actual% - Expected%) * ln(Actual% / Expected%) - 实现(用Evidently):
from evidently.metrics import DataDriftTable from evidently.report import Report report = Report(metrics=[DataDriftTable()]) report.run( reference_data=ref_df, # 上周数据 current_data=curr_df, # 今日数据 column_mapping={"numerical_features": ["feature_a", "feature_b"]} ) drift_result = report.as_dict() psi_values = drift_result["metrics"][0]["result"]["drift_by_columns"] for col, data in psi_values.items(): if data["drift_score"] > 0.15: alert(f"Feature {col} drift detected: PSI={data['drift_score']:.3f}")
3.6 日志结构化:用OpenTelemetry统一埋点
为什么?混用print()、logging.info()、structlog导致日志格式混乱,ELK里查个错误要写正则。
怎么做?
- 所有服务(
>from opentelemetry import trace from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter from opentelemetry.sdk.trace import TracerProvider from opentelemetry.sdk.trace.export import BatchSpanProcessor provider = TracerProvider() processor = BatchSpanProcessor(OTLPSpanExporter(endpoint="http://otel-collector:4318/v1/traces")) provider.add_span_processor(processor) trace.set_tracer_provider(provider) - 关键日志打点:
tracer = trace.get_tracer(__name__) with tracer.start_as_current_span("model.predict") as span: span.set_attribute("model.version", "v3.2") span.set_attribute("input.size", len(input_tensor)) result = model(input_tensor) span.set_attribute("output.score", float(result["score"]))
3.7 模型版本管理:Git LFS + S3双备份
为什么?只存S3,一旦误删不可逆;只存Git,大文件拖慢克隆。
怎么做?
.gitattributes配置:*.pt filter=lfs diff=lfs merge=lfs -text *.onnx filter=lfs diff=lfs merge=lfs -text- 推送时自动同步S3:
# Git hook: post-commit aws s3 sync ./models/ s3://my-bucket/models/ --exclude "*" --include "*.pt" --include "*.onnx"
3.8 安全加固:gRPC TLS双向认证
为什么?内部服务间调用若无认证,一个被攻陷的>cfssl gencert -initca ca-csr.json | cfssljson -bare ca cfssl gencert -ca=ca.pem -ca-key=ca-key.pem -config=ca-config.json -profile=server server-csr.json | cfssljson -bare server
model-runnergRPC Server启用TLS:server_credentials = grpc.ssl_server_credentials( [(open('server-key.pem', 'rb').read(), open('server.pem', 'rb').read())], root_certificates=open('ca.pem', 'rb').read(), require_client_auth=True ) server.add_secure_port('[::]:50051', server_credentials)3.9 部署编排:用Kustomize而非Helm管理环境差异
为什么?Helm模板复杂,values.yaml嵌套深,不同环境(dev/staging/prod)diff难读;Kustomize用patch机制,一眼看清prod比dev多了哪些配置。
怎么做?
- 目录结构:
k8s/ ├── base/ # 公共资源 │ ├── deployment.yaml │ └── service.yaml ├── dev/ │ ├── kustomization.yaml │ └── patch-env.yaml ├── prod/ │ ├── kustomization.yaml │ └── patch-resource.yaml # 增加CPU limit prod/kustomization.yaml:resources: - ../base patches: - patch-resource.yaml - path: patch-tls.yaml # 启用TLS
3.10 测试策略:分层测试金字塔
为什么?只测端到端,故障定位慢;只测单元,无法发现集成问题。
怎么做?
- 单元测试(70%):
model-runner的Tensor输入/输出校验,用pytest+torch.testing:def test_model_output_shape(): input = torch.randn(1, 100) output = model(input) assert output.shape == (1, 2) # 分类2类 - 集成测试(20%):
>@pytest.fixture(scope="session") def kafka_container(): return KafkaContainer("confluentinc/cp-kafka:7.3.0") - 混沌测试(10%):用Chaos Mesh随机kill
model-runnerPod,验证Orchestrator自动重试与降级。
3.11 降级策略:静态Fallback模型+规则兜底
为什么?模型服务宕机时,若只返回503,业务就断了。必须有“能用”的最低保障。
怎么做?
orchestrator内置Fallback:async function predictWithFallback(input: Input): Promise<Result> { try { return await callModelService(input); // 主路 } catch (e) { // 降级:用轻量规则引擎 return ruleEngine.evaluate(input); } }- 规则引擎预置兜底逻辑(如风控场景):
fallback_rules: - name: "high_risk_default" condition: "true" action: "block" confidence: 0.7 # 降级置信度,用于监控
3.12 成本监控:按模型实例粒度统计GPU小时
为什么?不细粒度计费,无法评估模型ROI;业务方总说“模型效果好”,但没人算过它每月烧掉多少云费用。
怎么做?
model-runner启动时上报元数据到Prometheus:from prometheus_client import Gauge gpu_hour_gauge = Gauge('model_gpu_hours_total', 'GPU hours consumed', ['model_name', 'version']) # 每分钟上报当前GPU使用率 gpu_hour_gauge.labels(model_name="fraud-detection", version="v3.2").set(gpu_util * 1/60)- Grafana面板:按模型名聚合,计算月度GPU小时成本($0.85/hour * GPU小时)。
这12个点,每个都来自真实战场。比如PSI检测,我们曾因忽略它,让一个推荐模型在促销季悄悄失效——用户点击率跌了40%,但A/B测试P95延迟没变,无人察觉,直到财务发现GMV异常。From Scratch不是炫技,是把每个环节的“意外”变成“预期”。
4. 实操全流程:从空目录到生产服务的72小时
现在,我们用一个具体案例走一遍完整流程:为电商APP上线“实时商品相似推荐”服务。目标:QPS 200,P95延迟 ≤150ms,支持每日增量训练。
4.1 Day 1:环境准备与骨架搭建(8小时)
目标:本地跑通最小可行服务链路(HTTP → Orchestrator → Mock Model)。
- 初始化项目:
mkdir ai-engineering-from-scratch && cd $_ git init echo "node_modules/" >> .gitignore - 搭建Orchestrator(TypeScript):
npm init -y npm install express @opentelemetry/api @opentelemetry/sdk-node npx tsc --init --rootDir src --outDir dist --strict - 编写
src/index.ts(极简版):import express from 'express'; const app = express(); app.use(express.json()); app.post('/recommend', async (req, res) => { // 模拟调用模型 const mockResult = { items: ["item_123", "item_456"], score: 0.92 }; res.json(mockResult); }); app.listen(3000, () => console.log('Orchestrator running on port 3000')); - 启动验证:
npx tsc && node dist/index.js curl -X POST http://localhost:3000/recommend -H "Content-Type: application/json" -d '{"user_id":"u123","item_id":"i789"}'
实操心得:第一天绝不碰真实模型!先确保HTTP服务、日志、监控埋点能跑通。我见过太多团队卡在PyTorch CUDA环境配置上,浪费三天——先让“壳”活起来,再往里填“肉”。
4.2 Day 2:数据管道与模型服务(12小时)
目标:接入真实数据源,model-runner能加载PyTorch模型并响应。
- 部署
>cargo new>let consumer = ConsumerConfig::new() .add_broker("localhost:9092") .create_with_context(KafkaContext).unwrap(); loop { let msg = consumer.recv_timeout(Duration::from_secs(1)).unwrap(); let record: UserEvent = AvroDeserializer::deserialize(&msg.payload()).unwrap(); // 发送到gRPC endpoint let mut client = ModelClient::connect("http://localhost:50051").await?; let request = PredictRequest { input: record.to_tensor() }; let response = client.predict(request).await?; println!("Score: {}", response.score); } - 构建
model-runner(Python):mkdir model-runner && cd $_ python3 -m venv venv && source venv/bin/activate pip install torch torchvision grpcio protobuf - 编写
server.py:import torch import grpc import model_pb2_grpc class ModelServicer(model_pb2_grpc.ModelServicer): def __init__(self): self.model = torch.jit.load("model.pt") def Predict(self, request, context): input_tensor = torch.tensor(request.input).reshape(1, -1) with torch.no_grad(): output = self.model(input_tensor) return model_pb2.PredictResponse(score=float(output[0][1])) - 生成gRPC stub(
model.proto):syntax = "proto3"; service Model { rpc Predict(PredictRequest) returns (PredictResponse); } message PredictRequest { repeated float input = 1; } message PredictResponse { float score = 1; } - 编译:
python -m grpc_tools.protoc -I. --python_out=. --grpc_python_out=. model.proto
注意:Day 2务必完成gRPC通信验证。用
grpcurl测试:grpcurl -plaintext -d '{"input":[0.1,0.2,0.3]}' localhost:50051 model.Model/Predict若失败,90%是Python/Go/Java生成的stub版本不一致,务必统一protoc版本。
4.3 Day 3:可观测性与自动化(16小时)
目标:所有服务暴露Prometheus指标,Grafana看板可用,CI/CD流水线就绪。
- 集成Prometheus:
orchestrator添加prom-client:import { collectDefaultMetrics, register } from 'prom-client'; collectDefaultMetrics(); // 自动收集Node.js指标 const httpRequestDuration = new Histogram({ name: 'http_request_duration_seconds', help: 'Duration of HTTP requests in seconds', labelNames: ['method', 'route', 'status_code'], buckets: [0.01, 0.05, 0.1, 0.2, 0.5, 1, 2, 5] }); app.use((req, res, next) => { const end = httpRequestDuration.startTimer({ method: req.method, route: req.route?.path || 'unknown' }); res.on('finish', () => end({ status_code: res.statusCode })); next(); }); app.get('/metrics', async (req, res) => { res.set('Content-Type', register.contentType); res.end(await register.metrics()); });
- 部署Prometheus+Grafana(Docker Compose):
# docker-compose.yml services: prometheus: image: prom/prometheus volumes: - ./prometheus.yml:/etc/prometheus/prometheus.yml grafana: image: grafana/grafana ports: - "3000:3000" - 编写CI/CD(GitHub Actions):
# .github/workflows/deploy.yml name: Deploy to Staging on: [push] jobs: deploy: runs-on: ubuntu-latest steps: - uses: actions/checkout@v3 - name: Build and Push Docker run: | docker build -t ${{ secrets.REGISTRY }}/orchestrator:${{ github.sha }} . docker push ${{ secrets.REGISTRY }}/orchestrator:${{ github.sha }} - name: Deploy to Kubernetes run: | kubectl set image deployment/orchestrator orchestrator=$SECRET_REGISTRY/orchestrator:${{ github.sha }}
4.4 Day 4:压力测试与调优(12小时)
目标:验证P95延迟≤150ms,找出瓶颈并优化。
- 用k6压测Orchestrator:
// script.js import http from 'k6/http'; export default function () { http.post('http://localhost:3000/recommend', JSON.stringify({ "user_id": "u" + __ENV.USER_ID, "item_id": "i" + __ENV.ITEM_ID }), { headers: { 'Content-Type': 'application/json' } }); } - 运行:
k6 run -e USER_ID=123 -e ITEM_ID=456 -d 5m -u 100 script.js - 分析结果:
- 若P95>150ms,先看Grafana:是Orchestrator CPU高?还是
model-runnergRPC延迟高? - 若
model-runner延迟高,用nvtop看GPU利用率;若<60%,说明模型未满载,需增加并发数。 - 若Orchestrator延迟高,用
pprof抓火焰图:常见瓶颈是JSON序列化(换simd-json)或日志格式化(用pino替代winston)。
- 若P95>150ms,先看Grafana:是Orchestrator CPU高?还是
实操心得:压力测试必须模拟真实数据分布。我们曾用均匀随机ID压测,P95达标;上线后真实用户ID有长尾(热门商品被刷),缓存命中率暴跌,延迟翻倍。解决方案:用线上一周访问日志抽样生成压测数据集。
4.5 Day 5:安全加固与合规(8小时)
目标:通过内部安全审计,满足GDPR数据最小化原则。
- 实施TLS双向认证(见3.8)。
- 数据脱敏:
>use sha2::{Sha256, Digest}; let salted = format!("{}{}", user_id, "my_secret_salt"); let mut hasher = Sha256::new(); hasher.update(salted); let hash = hasher.finalize(); - 添加数据保留策略:所有原始事件在S3中自动过期(生命周期规则设为30天)。
- 生成SOC2合规报告:用
checkov扫描IaC(Kustomize YAML),用bandit扫描Python代码。
4.6 Day 6:上线与监控(8小时)
目标:灰度发布,建立基线,设置告警。
- 配置灰度(Orchestrator UI):对1%用户启用新服务,99%走旧API。
- 建立监控基线(Grafana):
- P95延迟基线:首小时平均值±10%
- 错误率基线:0.1%
- GPU利用率基线:45%
- 设置PagerDuty告警:
rate(http_request_duration_seconds_bucket{le="0.15"}[5m]) / rate(http_request_duration_seconds_count[5m]) < 0.95(P95超150ms)sum(rate(grpc_server_handled_total{job="model-runner"}[5m])) by (grpc_code) > 0 and sum(rate(grpc_server_handled_total{grpc_code="Unknown"}[5m])) by (grpc_code) > 10(未知错误突增)
4.7 Day 7:反馈闭环与迭代(8小时)
目标:确认反馈数据流入,首个bad case分析完成。
- 验证Delta Lake数据写入:
-- Spark SQL SELECT COUNT(*) FROM delta.`s3://bucket/predictions/` WHERE date = '2023-10-01'; - 运行漂移检测Job:
spark-submit \ --class com.evidently.DriftDetector \ --master yarn \ drift-detector.jar \ --ref-path s3://bucket/ref-data/ \ --curr-path s3://bucket/today-data/ \ --output s3://bucket/drift-report/ - 分析首个bad case:从审计日志中提取10个被人工覆盖的预测,对比模型输出与真实结果,定位特征偏差(如新上架商品缺少历史销量特征)。
72小时不是神话,是我们团队的标准交付节奏。关键不是速度,而是