ARTICLE DETAIL

资讯详情

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

从零搭建AI工程体系:数据管道、特征存储与模型部署实战

从零搭建AI工程体系:数据管道、特征存储与模型部署实战 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, conengine) # 计算数据指纹 fingerprint hashlib.md5( pd.util.hash_pandas_object(df).values.tobytes() ).hexdigest() # 存储快照 path fs3://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__( nameuser_purchase_count_7d, dtypeint64, dependencies[user_id, event_time, event_type], window7d, default_value0, ownerdata-team ) def compute(self, df: pd.DataFrame) - pd.Series: cutoff df[event_time].max() - pd.Timedelta(days7) 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:ibatch_size] pipe redis_client.pipeline() for _, row in batch.iterrows(): key ffeature:{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, binsbreakpoints)[0] / len(expected) actual_counts np.histogram(actual, binsbreakpoints)[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( dsnpostgresql://..., min_size10, max_size50 ) app.state.redis await aioredis.create_redis_pool( redis://..., minsize10, maxsize50 ) yield # 关闭时释放连接 await app.state.db_pool.close() app.state.redis.close() app FastAPI(lifespanlifespan) 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%的流量走旧模型作为对照组。这样即使新模型有问题你也能通过对比及时发现。这个策略叫“影子流量”实现成本很低但价值极高。我靠这个策略避免了好几次重大事故。
返回列表