ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

AI工程从零构建:数据、特征、模型、服务、监控五大核心模块实战

AI工程从零构建:数据、特征、模型、服务、监控五大核心模块实战 1. 这不是调包是亲手搭起AI工程的骨架“AI Engineering from Scratch”——看到这个标题很多人第一反应是又要学Python又要装CUDA又要配环境别急先放下这些预设。我带过二十多个从零起步的AI工程落地项目最常听到的抱怨不是“模型不会训”而是“代码跑通了但上线就崩”“本地效果好一上生产就掉点”“团队协作时连数据版本都对不上”。这根本不是算法问题是工程能力断层。AI Engineering from Scratch说白了就是拒绝把“能跑通”当交付而是用软件工程的标准重新定义AI项目的生命周期从数据怎么存、特征怎么管、模型怎么验、服务怎么扩、监控怎么设到故障怎么回滚——每一环都得自己搭、自己测、自己压、自己盯。它不教你怎么调参而是教你怎么让一个模型在真实业务里活过三个月它不讲Transformer原理但会告诉你为什么feature store必须带schema校验为什么模型注册表要强制绑定训练数据哈希为什么API响应延迟超过800ms就必须触发降级开关。关键词“ai-engineering”和“from-scratch”不是修饰词是底线所有基础设施、所有流程规范、所有可观测性组件全部手写或深度定制不依赖黑盒平台不跳过任何中间环节。适合三类人刚转行想真正理解AI系统全貌的工程师技术负责人需要搭建可审计、可复现、可交接的AI基建还有那些被“MLOps平台一键部署”坑过、发现线上问题根本查不到源头的实战派。这不是速成课这是给AI系统立规矩。2. 为什么非得从零开始绕不开的五个工程断点2.1 数据管道不是ETL是可信度流水线很多团队把数据准备当成“前置步骤”训完模型才回头补数据质量报告。结果呢线上推理时突然发现某类用户ID字段里混进了空格和emoji特征计算直接NaN服务返回500。从零构建AI工程第一刀就砍在数据管道上。它不能只是“把csv读进来、清洗、喂给模型”而是一条带校验、带版本、带血缘的可信度流水线。我做过一个风控模型原始数据来自三个业务库字段命名混乱user_id、uid、customer_no时间戳格式不统一ISO、Unix、毫秒字符串缺失值填充策略各不相同。如果用现成的Airflow DAG简单串起来上线后才发现A库某天凌晨2点因数据库锁表延迟3小时同步导致当天特征计算用了前一日快照误杀率飙升17%。从零设计时我们强制每个数据源接入点必须声明SLA最大延迟容忍、schema字段类型非空约束枚举值范围、以及变更通知机制如DDL变更自动触发pipeline冻结。清洗阶段不是写一堆pandas .fillna()而是定义原子化校验器NullChecker、RangeValidator、ConsistencyCrossRef比如订单金额必须大于等于商品单价×数量。每个校验失败都会生成结构化告警包含样本ID、失败字段、偏离值、影响行数并自动阻断下游任务。实操中我们用Pydantic v2定义schema用Great Expectations做校验规则编排但核心不是工具而是把“数据可信”变成可量化、可追踪、可追责的工程契约。绕开这点后面所有模型都是沙上筑塔。2.2 特征管理不是缓存是状态机驱动的契约中心“特征工程”常被简化为“写SQL提取字段”但真实场景里特征是动态的、有生命周期的、需协同的。比如一个电商推荐系统的“用户7日加购频次”上线后运营要求增加“排除促销商品”算法同学想加入“品类偏好衰减因子”风控团队又提出“需隔离羊毛党行为”。如果特征逻辑散落在各个notebook和SQL脚本里改一处全链路崩。从零构建时我们把Feature Store做成一个状态机驱动的契约中心每个特征注册时必须声明version、owner、update_frequency、serving_latency_sla、backfill_window以及最关键的dependency_graph依赖哪些上游表、哪些其他特征。例如user_7d_cart_count_v2明确依赖raw_events表和item_category_mapping_v1特征任何上游变更都会触发自动影响分析。我们不用Flink实时计算而是用Dask分布式批处理Redis缓存但重点在于所有特征计算代码必须通过单元测试mock数据输入验证输出schema和数值边界所有特征上线前必须通过A/B测试流量切分验证对比v1和v2在相同样本上的分布偏移KS值0.01。最深的坑是时间旅行问题线上服务需要获取“用户在T时刻的特征”但特征计算本身有延迟。我们采用“事件时间处理时间双时间戳”机制特征存储中每条记录带event_time用户行为发生时间和ingestion_time特征入库时间服务层根据请求时间戳自动选择最近可用快照。没这套机制特征漂移就是定时炸弹。2.3 模型注册不是存文件是带上下文的可追溯实体把.pkl或.onnx文件扔进S3桶叫“模型存储”不叫“模型注册”。从零构建的模型注册表本质是一个带完整上下文的可追溯实体。它必须强制关联训练代码commit hash、训练数据集版本含data catalog ID和sample hash、超参配置JSON Schema校验、评估指标精确到每个子集的F1/Recall/AUC、以及人工审核记录谁、何时、基于什么证据批准上线。我们曾遇到一个NLP模型在测试集上准确率92%但上线后发现对长尾行业词如“量子计算芯片封装工艺”完全失效。回溯发现训练数据中该类样本仅占0.03%且标注质量差但评估报告只汇报了macro-average掩盖了问题。从零设计时注册表强制要求多维度评估报告按行业、按文本长度、按实体密度分组统计并生成可视化分布图。模型加载时服务端会校验当前运行环境Python版本、torch版本、CUDA驱动是否与注册时声明的environment_spec兼容不匹配则拒绝加载并告警。更关键的是版本策略我们不用简单的v1/v2而是采用model_name-YYYYMMDD-git_short_hash格式确保每次变更都有唯一、可定位的标识。一次线上事故中运维同事5分钟内就定位到是fraud_model-20240315-a7f2b1c版本引入了新特征缩放逻辑回滚到fraud_model-20240310-8d4e92f即恢复全程无需翻代码、无需问算法同学。2.4 服务部署不是起个Flask是弹性与弹性的博弈“用FastAPI跑个predict endpoint”只是起点。从零构建的服务层核心矛盾是弹性伸缩与推理稳定性的博弈。GPU资源昂贵但突发流量可能瞬间打满显存CPU服务便宜但复杂模型推理延迟波动大。我们放弃Kubernetes HPA的默认CPU/Memory指标自研基于QPS和P99延迟的混合伸缩策略。服务启动时每个worker进程主动上报自身负载能力基线warmup阶段用固定样本测100次取P99然后每30秒向中央协调器上报实时指标当前QPS、P99延迟、GPU显存占用率、请求队列长度。协调器根据预设SLA如P99500ms错误率0.1%动态调整副本数。更关键的是熔断与降级当单实例P99连续5次800ms自动触发熔断将流量切换至轻量级fallback模型如LR规则引擎当整体错误率1%启动渐进式限流令牌桶漏桶双机制。实操中我们用Prometheus采集指标用Grafana看板实时监控但真正的工程价值在于所有熔断阈值、降级策略、扩容步长都写死在服务配置中而非运维手动干预。一次大促期间主模型因数据分布突变导致延迟飙升系统在23秒内完成熔断降级扩容用户无感知而传统方案依赖人工告警-登录-排查-操作平均耗时6分钟以上。2.5 监控告警不是看曲线是定义业务健康的语言“GPU显存使用率90%告警”毫无意义。从零构建的监控体系必须用业务语言定义健康。我们把监控拆成三层基础设施层GPU温度、NVLink带宽、服务层API成功率、P99延迟、特征计算耗时、业务层模型预测置信度分布偏移、关键特征值域漂移、线上A/B测试指标衰减。其中业务层监控最难也最重要。例如风控模型我们监控high_risk_prediction_rate高风险预测占比的7日滑动标准差若0.05则告警——这比单纯看准确率下降更能提前发现数据漂移。再如推荐系统监控top_k_diversity_score推荐列表品类多样性若连续3小时低于阈值则触发特征新鲜度检查。所有告警规则必须关联根因预案high_risk_rate_std_alert自动触发“检查近24小时用户地域分布变化”和“拉取最新样本重跑特征漂移检测”。我们不用ELK堆日志而是用OpenTelemetry统一埋点所有trace span都注入model_version、feature_version、request_id标签确保一次异常请求能10秒内定位到具体模型、特征、数据批次。没有这套以业务结果为导向的监控所谓“可观测性”只是仪表盘上漂亮的曲线。3. 核心模块手把手实现不跳过一行关键代码3.1 可信数据管道Schema驱动的校验引擎数据管道的基石是schema定义与校验。我们不用Apache Avro或Protobuf而是用Pydantic v2定义轻量级、可执行的schema因为它支持运行时校验和自定义validator。以下是一个典型用户行为数据schemafrom pydantic import BaseModel, validator, Field from typing import Optional, List, Dict, Any import re class UserEventSchema(BaseModel): event_id: str Field(..., min_length10, max_length32) user_id: str Field(..., regexr^[a-zA-Z0-9_]{8,32}$) # 强制格式 event_type: str Field(..., pattern^(click|view|purchase|cart_add)$) timestamp_ms: int Field(..., ge1609459200000, le2524608000000) # 2021-2050 item_id: Optional[str] None category_path: Optional[str] None price_cents: Optional[int] Field(None, ge0, le100000000) # 最大100万人民币 validator(category_path) def validate_category_path(cls, v): if v is not None and not re.match(r^[a-zA-Z0-9_](\/[a-zA-Z0-9_])*$, v): raise ValueError(category_path must be slash-separated alphanumeric segments) return v validator(price_cents) def validate_price_cents(cls, v, values): if v is not None and event_type in values and values[event_type] purchase: if v 0: raise ValueError(purchase event must have positive price_cents) return v校验引擎核心是DataValidator类它接收原始字典列表批量校验并生成结构化报告from collections import defaultdict import json class DataValidator: def __init__(self, schema_class: type[BaseModel]): self.schema_class schema_class self.errors defaultdict(list) # {error_type: [(row_idx, error_msg), ...]} def validate_batch(self, records: List[Dict[str, Any]]) - bool: 返回True表示全部通过False表示存在错误 self.errors.clear() for idx, record in enumerate(records): try: # Pydantic自动校验并转换类型 self.schema_class(**record) except Exception as e: # 提取Pydantic错误信息标准化为{field, error_type, message} if hasattr(e, errors): for err in e.errors(): error_key f{err[loc][0]}_{err[type]} if len(err[loc]) 0 else root_validation self.errors[error_key].append((idx, err[msg])) else: self.errors[unknown_error].append((idx, str(e))) return len(self.errors) 0 def get_report(self) - Dict[str, Any]: 生成JSON序列化报告供告警和审计 total_records sum(len(v) for v in self.errors.values()) return { total_errors: total_records, error_summary: {k: len(v) for k, v in self.errors.items()}, sample_errors: [ {row_index: idx, error_type: error_type.split(_)[0], message: msg} for error_type, errors in list(self.errors.items())[:3] for idx, msg in errors[:2] ] } # 使用示例 validator DataValidator(UserEventSchema) raw_data [ {event_id: evt_123, user_id: u123, event_type: click, timestamp_ms: 1710000000000}, {event_id: evt_456, user_id: u!#, event_type: purchase, timestamp_ms: 1710000000000, price_cents: -100} # 错误user_id非法price_cents负数 ] is_valid validator.validate_batch(raw_data) print(fValid: {is_valid}) # False print(json.dumps(validator.get_report(), indent2))提示实际生产中validate_batch会集成到Spark或Dask任务中每个分区独立校验错误样本自动写入quarantine目录并触发告警。关键不是校验本身而是让错误可定位、可归因、可追溯——每个错误都绑定原始行号和字段名避免“数据有问题”这种模糊反馈。3.2 特征注册中心带依赖图谱的版本化仓库特征注册的核心是解决“谁在用、谁在改、改了影响谁”的问题。我们用SQLite做轻量级元数据存储避免引入复杂DB但重点在依赖关系建模。每个特征注册为FeatureDef对象from dataclasses import dataclass, field from datetime import datetime from typing import List, Optional, Dict, Any import hashlib import json dataclass class FeatureDef: name: str version: str # 格式v1.2.0 或 20240315 owner: str # 邮箱或团队名 description: str dependencies: List[str] # [raw_events_v2, user_profile_v1] update_frequency: str # hourly, daily, realtime serving_latency_sla_ms: int backfill_window_days: int code_hash: str # 计算特征计算代码的sha256 created_at: datetime field(default_factorydatetime.now) updated_at: datetime field(default_factorydatetime.now) def to_dict(self) - Dict[str, Any]: d self.__dict__.copy() d[created_at] self.created_at.isoformat() d[updated_at] self.updated_at.isoformat() return d def calculate_code_hash(self, code_content: str) - str: 计算特征计算代码的hash用于检测逻辑变更 return hashlib.sha256(code_content.encode()).hexdigest()[:12] # 特征注册中心类 class FeatureRegistry: def __init__(self, db_path: str feature_registry.db): self.db_path db_path self._init_db() def _init_db(self): import sqlite3 conn sqlite3.connect(self.db_path) conn.execute( CREATE TABLE IF NOT EXISTS features ( id INTEGER PRIMARY KEY AUTOINCREMENT, name TEXT NOT NULL, version TEXT NOT NULL, owner TEXT NOT NULL, description TEXT, dependencies TEXT, -- JSON array update_frequency TEXT, serving_latency_sla_ms INTEGER, backfill_window_days INTEGER, code_hash TEXT, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, UNIQUE(name, version) ) ) conn.close() def register_feature(self, feature_def: FeatureDef, code_content: str): 注册新特征自动计算code_hash feature_def.code_hash feature_def.calculate_code_hash(code_content) import sqlite3 conn sqlite3.connect(self.db_path) conn.execute( INSERT INTO features (name, version, owner, description, dependencies, update_frequency, serving_latency_sla_ms, backfill_window_days, code_hash) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) , ( feature_def.name, feature_def.version, feature_def.owner, feature_def.description, json.dumps(feature_def.dependencies), feature_def.update_frequency, feature_def.serving_latency_sla_ms, feature_def.backfill_window_days, feature_def.code_hash )) conn.commit() conn.close() def get_feature_by_name(self, name: str, version: str None) - Optional[FeatureDef]: 获取特征定义支持最新版本或指定版本 import sqlite3 conn sqlite3.connect(self.db_path) if version is None: cursor conn.execute( SELECT * FROM features WHERE name ? ORDER BY created_at DESC LIMIT 1 , (name,)) else: cursor conn.execute( SELECT * FROM features WHERE name ? AND version ? , (name, version)) row cursor.fetchone() if row: # 解析dependencies JSON deps json.loads(row[5]) if row[5] else [] return FeatureDef( namerow[1], versionrow[2], ownerrow[3], descriptionrow[4], dependenciesdeps, update_frequencyrow[6], serving_latency_sla_msrow[7], backfill_window_daysrow[8], code_hashrow[9], created_atdatetime.fromisoformat(row[10]), updated_atdatetime.fromisoformat(row[11]) ) return None def get_dependents(self, feature_name: str) - List[str]: 查询依赖指定特征的所有特征反向依赖 import sqlite3 conn sqlite3.connect(self.db_path) cursor conn.execute( SELECT name, version FROM features WHERE dependencies LIKE ? OR dependencies LIKE ? , (f%{feature_name}%, f%{feature_name}%)) return [f{row[0]}-{row[1]} for row in cursor.fetchall()] # 使用示例 registry FeatureRegistry() feature FeatureDef( nameuser_7d_cart_count, versionv2.1.0, ownerrecommender-teamcompany.com, description用户过去7天加购商品次数排除促销商品, dependencies[raw_events_v3, promotion_flag_v1], update_frequencydaily, serving_latency_sla_ms200, backfill_window_days30 ) registry.register_feature(feature, def compute_user_cart_count(events_df, promo_df): # 实际计算逻辑 pass )注意get_dependents方法是关键。当promotion_flag_v1更新时系统自动调用此方法找出所有依赖它的特征如user_7d_cart_count_v2.1.0、user_conversion_rate_v1.0.0并触发它们的回归测试和重新计算。这解决了特征变更的连锁反应问题避免“改一个小特征崩掉十个模型”。3.3 模型注册表带环境约束的可追溯实体模型注册表必须强制绑定运行时环境否则“本地能跑”和“线上能跑”永远是两回事。我们用JSON Schema定义环境约束并在模型加载时严格校验import jsonschema import platform import sys import torch # 环境约束Schema ENV_SCHEMA { type: object, properties: { python_version: {type: string, pattern: r^\d\.\d\.\d$}, torch_version: {type: string, pattern: r^\d\.\d\.\d$}, cuda_version: {type: string, pattern: r^\d\.\d$}, os_platform: {enum: [linux, darwin, win32]}, gpu_count: {type: integer, minimum: 0} }, required: [python_version, torch_version, os_platform] } class ModelRegistry: def __init__(self, registry_dir: str model_registry): self.registry_dir registry_dir import os os.makedirs(registry_dir, exist_okTrue) def register_model(self, model_name: str, version: str, model_file: str, metadata: Dict[str, Any]): 注册模型metadata必须包含env_spec import shutil import json from datetime import datetime # 校验env_spec try: jsonschema.validate(instancemetadata.get(env_spec, {}), schemaENV_SCHEMA) except jsonschema.ValidationError as e: raise ValueError(fInvalid env_spec: {e.message}) # 构建版本目录 version_dir f{self.registry_dir}/{model_name}/{version} import os os.makedirs(version_dir, exist_okTrue) # 复制模型文件 shutil.copy(model_file, f{version_dir}/model.bin) # 保存元数据 metadata_full { model_name: model_name, version: version, registered_at: datetime.now().isoformat(), env_spec: metadata[env_spec], training_data_hash: metadata.get(training_data_hash), eval_metrics: metadata.get(eval_metrics, {}), code_commit: metadata.get(code_commit), owner: metadata.get(owner) } with open(f{version_dir}/metadata.json, w) as f: json.dump(metadata_full, f, indent2) def load_model(self, model_name: str, version: str) - Any: 加载模型先校验环境兼容性 import json import torch version_dir f{self.registry_dir}/{model_name}/{version} if not os.path.exists(version_dir): raise FileNotFoundError(fModel {model_name}-{version} not found) # 读取元数据 with open(f{version_dir}/metadata.json) as f: meta json.load(f) # 校验环境 env_spec meta[env_spec] current_env { python_version: ..join(map(str, sys.version_info[:3])), torch_version: torch.__version__, os_platform: sys.platform, gpu_count: torch.cuda.device_count() if torch.cuda.is_available() else 0 } # 粗粒度兼容性检查版本号前缀匹配 if not self._is_compatible(env_spec, current_env): raise RuntimeError( fEnvironment mismatch: registered {env_spec} vs current {current_env} ) # 加载模型 model_path f{version_dir}/model.bin if model_path.endswith(.bin): # 假设是PyTorch模型 model torch.load(model_path, map_locationcpu) model.eval() return model else: raise NotImplementedError(Only .bin models supported) def _is_compatible(self, required: Dict[str, str], current: Dict[str, Any]) - bool: 检查当前环境是否满足required约束 for key, req_val in required.items(): if key not in current: return False cur_val current[key] if isinstance(cur_val, str) and isinstance(req_val, str): # 版本号兼容req_val1.12.0 兼容 current1.12.1 if key in [python_version, torch_version]: req_parts req_val.split(.)[:2] # 取主版本和次版本 cur_parts cur_val.split(.)[:2] if req_parts ! cur_parts: return False elif key cuda_version: # CUDA版本需完全匹配 if req_val ! cur_val: return False elif key os_platform: if req_val ! cur_val: return False elif key gpu_count: if cur_val req_val: # 要求至少req_val个GPU return False return True # 使用示例 registry ModelRegistry() registry.register_model( model_namefraud_detector, version20240315-a7f2b1c, model_filemodels/fraud_v1.bin, metadata{ env_spec: { python_version: 3.9.16, torch_version: 1.13.1, cuda_version: 11.7, os_platform: linux, gpu_count: 1 }, training_data_hash: sha256:abc123..., eval_metrics: {auc: 0.92, f1: 0.85}, code_commit: a7f2b1c, owner: risk-teamcompany.com } ) # 加载时自动校验 try: model registry.load_model(fraud_detector, 20240315-a7f2b1c) except RuntimeError as e: print(fLoad failed: {e}) # 环境不匹配时抛出明确错误实操心得_is_compatible方法中的版本兼容策略是经验之谈。PyTorch主次版本如1.12.x通常ABI兼容但补丁版本x可能有细微差异所以只校验前两位CUDA版本必须严格匹配因为驱动和runtime的二进制接口不向前兼容GPU数量是硬性要求少于注册值会导致OOM。这套机制让模型部署从“祈祷能跑”变成“确定能跑”。3.4 服务治理基于延迟的混合伸缩控制器服务伸缩不能只看CPU必须结合业务指标。我们实现了一个轻量级控制器每30秒采集指标并决策import time import threading import requests from typing import Dict, Any, List import logging class ScalingController: def __init__(self, service_url: str, target_p99_ms: int 500, min_replicas: int 1, max_replicas: int 20): self.service_url service_url self.target_p99_ms target_p99_ms self.min_replicas min_replicas self.max_replicas max_replicas self.current_replicas min_replicas self.metrics_history [] # 存储最近10次指标 self.lock threading.Lock() self.logger logging.getLogger(__name__) def collect_metrics(self) - Dict[str, Any]: 采集服务指标QPS, P99延迟, 错误率, GPU显存 try: # 调用服务健康端点需服务暴露/metrics resp requests.get(f{self.service_url}/metrics, timeout5) if resp.status_code 200: metrics resp.json() # 示例metrics: {qps: 120.5, p99_ms: 420.3, error_rate: 0.002, gpu_mem_percent: 75.2} return metrics else: self.logger.warning(fMetrics endpoint returned {resp.status_code}) return {} except Exception as e: self.logger.error(fFailed to collect metrics: {e}) return {} def calculate_target_replicas(self, metrics: Dict[str, Any]) - int: 基于指标计算目标副本数 if not metrics or qps not in metrics or p99_ms not in metrics: return self.current_replicas qps metrics[qps] p99 metrics[p99_ms] error_rate metrics.get(error_rate, 0.0) # 核心策略优先保障延迟SLA if p99 self.target_p99_ms * 1.5: # 严重超标激进扩容 scale_factor min(2.0, p99 / self.target_p99_ms) target int(self.current_replicas * scale_factor) elif p99 self.target_p99_ms * 1.2: # 轻微超标温和扩容 target self.current_replicas 1 elif p99 self.target_p99_ms * 0.8 and qps self.current_replicas * 50: # 低负载且达标缩容 target max(self.min_replicas, self.current_replicas - 1) else: # 达标维持 target self.current_replicas # 错误率兜底1%立即扩容 if error_rate 0.01: target min(self.max_replicas, self.current_replicas 2) return max(self.min_replicas, min(self.max_replicas, target)) def adjust_replicas(self, target_replicas: int): 调用K8s API或服务管理API调整副本数 # 这里模拟调用K8s patch if target_replicas ! self.current_replicas: self.logger.info(fScaling from {self.current_replicas} to {target_replicas} replicas) # 实际代码调用K8s API patch deployment # requests.patch(https://k8s/api/v1/namespaces/default/deployments/my-service, ...) self.current_replicas target_replicas def run_loop(self): 主循环 while True: try: metrics self.collect_metrics() if metrics: self.metrics_history.append(metrics) if len(self.metrics_history) 10: self.metrics_history.pop(0) target self.calculate_target_replicas(metrics) self.adjust_replicas(target) time.sleep(30) # 每30秒执行一次 except Exception as e: self.logger.error(fScaling loop error: {e}) time.sleep(30) # 启动控制器 controller ScalingController( service_urlhttp://my-ai-service:8000, target_p99_ms500, min_replicas2, max_replicas10 ) threading.Thread(targetcontroller.run_loop, daemonTrue).start()关键细节calculate_target_replicas函数体现了工程权衡。它不追求理论最优而是设定清晰的业务规则P99超标1.5倍以上才激进扩容避免毛刺误判达标且QPS低于每副本50 QPS才缩容防止频繁抖动错误率1%立即行动业务不可接受。这种策略比纯数学公式更可靠因为AI服务的负载模式高度非线性。4. 真实踩坑记录那些文档里绝不会写的教训4.1 “数据版本”陷阱你以为的同一份数据其实早已分裂最隐蔽的坑是数据版本漂移。我们曾有一个模型离线评估AUC 0.92上线后跌到0.78。排查三天最终发现训练时用的user_features.parquet是2024-03-10生成的而线上服务读取的同名文件是运维同学上周清理磁盘时从备份恢复的2024-02-20旧版。两者schema相同但user_age_group字段的枚举值从[0-18,19-35,36-50,51]变成了[under_18,18_35,35_50,over_50]模型把under_18当成了未知类别全部归为默认类。教训所有数据文件必须带不可篡改的版本标识。我们后来强制要求1Parquet文件名包含v{unix_timestamp}2文件metadata中写入data_hash整个文件内容的sha2563服务启动时校验data_hash与注册表中记录的一致。再没出现过此类问题。4.2 “特征一致性”幻觉训练和服务的特征计算根本不是同一段代码算法同学在Jupyter里写特征计算导出为feature_utils.py然后运维打包进服务镜像。看似一致实则危险。一次更新中算法同学修复了compute_user_ltv()函数里的一个除零bug但忘了同步feature_utils.py只更新了notebook。结果训练用新逻辑服务用旧逻辑特征值偏差达300%。教训特征计算代码必须单一源且由CI/CD自动同步。我们现在要求所有特征计算函数必须定义在features/目录下训练脚本和服务代码都import同一份CI流程中任何修改都触发全链路测试包括特征值一致性校验用相同输入比对训练和在线计算的输出diff。4.3 “模型热更新”迷思无缝切换背后是灾难温床追求“不重启服务更新模型”很诱人但实践中充满陷阱。我们曾用torch.jit.load()动态加载新模型结果发现1新模型权重加载后旧模型的GPU显存未释放内存持续增长2多线程并发加载时模型参数被意外覆盖3加载失败时服务无降级直接500。教训**模型更新必须
返回列表