在企业里做AI平台最难和别人解释清楚的往往不是模型而是数据。很多人下意识觉得AI平台就是GPU集群加上模型仓库真正跑起来才发现卡脖子的十有八九是数据处理这层。作为AI应用架构师我这两年最大的体会是数据处理不是模型训练前的准备工作而是整个AI平台的地基工程。它决定了你能多快拿到一份可信的训练集决定了特征上线之后会不会悄悄漂移决定了你从“能跑demo”到“能稳定交付模型”之间到底隔多远。这个主题适合所有正在搭中台、算法平台、训推平台的团队。我自己接手过的项目里数据散落在Hive、Kafka、业务库里标签口径对不上特征训练和上线不一致模型复现不了这些几乎成了标配问题。下面我按一个企业AI平台从零到落地的推进顺序把架构设计层面和数据处理实操层面的思路、选型、参数和踩坑一次性讲透。1. 先想清楚AI平台数据处理和传统ETL到底差在哪1.1 传统数据处理与AI数据工程的区别企业里做数据仓库的团队核心任务是ETL把数据从源系统搬运、清洗、聚合到数仓模型里最终产出给BI报表看。这套体系里数据模型是老早就定义好的维度加事实表结构消费者是人和看板。但AI平台的数据处理服务对象不是人和报表而是训练任务、评估任务和在线推理。这个本质的区别决定了架构上所有关键取舍。传统数仓看重“快、准、全”AI平台的数据更看重“可复现、可回放、可比较”。模型训练说到底是对历史数据的统计拟合如果数据管道在下一次跑的时候悄悄改了过滤逻辑哪怕只是改了一个时间窗口模型指标就会变。你拿着一个结果去复盘如果发现数据版本对不上整个团队的时间就耗在排查数据差异上而不是优化模型。所以架构师在考虑AI平台的数据处理时第一件事不是选框架不是搭集群而是把这些问题回答清楚数据怎么进来、怎么变干净、怎么被版本化、特征怎么保证训练和上线一致、数据万一坏了怎么回滚。这几件事在纸上画明白工具选型自然水到渠成。1.2 数据处理层在AI平台架构里的位置AI平台自底向上大致分五层基础设施层GPU/CPU、存储、网络、数据层采集、清洗、特征、版本管理、开发训练层实验管理、训练框架、模型仓库、服务层在线推理、模型评估、监控、上层业务应用。在这五层里数据层是承重墙。模型训练出了问题十有八九要先查数据算法工程师第一反应永远是“数据是不是又出问题了”这不是偏见而是数据层确实最容易埋雷。同时要注意数据处理不是一次性任务。它有两条明确的时间线离线那条线每天或每周批量产出训练集在线那条线按毫秒级产出特征供推理服务消费。两条线必须共享同一套元数据和特征定义否则就会出现典型的“训练时用特征A、上线时用特征B”事故这也是后文要重点讲的特征平台为什么重要。提示如果团队还在纠结要不要上Feature Store可以先画一张“特征从离线到在线”的流向图看看特征定义在两条链路上是否有重复实现。有重复就有隐患。2. 数据接入与清洗实操源头做对一半后面少走一半弯路2.1 接入层设计离线批量与线上流式两条腿走路企业里的数据源头五花八门业务数据库MySQL、Oracle、埋点日志、文件服务器、第三方接口。接入层要做的就是把这些来源统一汇入一层可被AI平台消费的存储。我建议先按“离线段”和“流式段”两条路线分别设计。离线段用批量同步从业务库定时拉全量或增量或者在HDFS上落地原始日志。工具上DataX、Sqoop、自研调度脚本都行关键不在工具而在三个细节全量和增量数据怎么配合、主键怎么去重、时区怎么统一。增量同步漏了数据全量重跑又覆盖了业务方手工修复的数据这类问题在接入层出现得最多。流式段用Kafka接埋点和CDC变更事件再用Flink做实时清洗和拼接。流式不是万能的但业务对T0有要求的时候它基本上是标配。这里要给个反面提醒很多团队为了追新技术把根本不需要实时的数据也强行上流式。流式会引入状态管理、水位线、幂等消费一大堆复杂度如果业务只是每天看一次报表离线完全够用没必要一开始就全面流式化。2.2 清洗环节缺失值、异常值与数据倾斜的处理细节数据清洗是数据处理里最重、最脏的活。我观察到一个规律真正影响模型效果的问题往往不是算法调参而是清洗阶段的细节。拿缺失值来说传统BI遇到缺失可以直接留空模型可不行。用DataFrame处理缺失值时第一步要区分三种情况字段缺失、空字符串、业务自定义的“无意义值”。这三种的处理方式完全不同字段缺失可以填充或删除空字符串可能是真实业务结果-999这种占位符如果当成真实数值特征分布会直接歪掉。我自己就在项目里踩过这个坑。一张用户行为表里时间戳字段偶尔是空字符串清洗脚本里写了fillna默认值但因为空字符串不是NaN填充根本没生效后面下游转换类型时又把空串变成了一个极早的时间戳模型上线后特定人群预测明显偏差。排查花了整整两天修复就一行先把空字符串统一映射成NaN再统一填充。现在我的清洗脚本一律强制先做类型标准化把、’NULL‘、’null‘这些全部归一化成np.nan再走填充逻辑。异常值处理也是类似。看到离群值不能上来就删先判断是传感器故障、录入错误还是真实业务极端值。我会先用分位数和IQR把候选异常值筛出来再结合业务背景决定保留还是剔除。有些场景里异常值是模型要学习的关键信号比如风控里的极端交易直接删掉反而是帮倒忙。数据倾斜是清洗阶段另一个高频问题尤其是用户特征按user_id做group by的场景头部用户占了大头数据量单容器直接OOM整个任务被拖垮。常用解法是加盐打散热点key后做二次汇聚或者对超大类单独拆分配置。这个优化方向正确的情况下任务耗时能从小时级降到十几分钟。2.3 数据剖面校验先发现问题别等训练炸了再排查清洗任务写完不要直接丢给训练。先跑一遍数据剖面Data Profiling输出每一列的均值、方差、分位数、缺失率、类别分布目的是在模型训练前发现“这波数据和上波长得不一样”。我见过太多情况是模型Loss突然变大排查半天发现是上游某张表新增了一个字段或者某个枚举值的分布发生了大变化。数据剖面的结果也应该按日分区保留它就是数据质量监控的天然基线后文第5部分会详细展开。这里先给一个最朴素的经验请务必让每个清洗任务在结束日志里输出一张“数据摘要表”这样哪怕后面发现问题也能快速定位是哪一天的哪份数据出了问题。3. 流式处理与批流一体实时模型不是加个Kafka就行3.1 从Lambda到Kappa为什么大家都在谈批流一体不少团队一说“要实时”架构图里加一个Kafka就以为完事了。真正落地时最棘手的问题是一致性同一份数据离线算出来50万条实时算出来49万条一对比报表就穿帮。早期典型的Lambda架构用两条独立路径分别处理批和流最后合并结果。这种方案能跑但代价是同一套业务逻辑维护两套代码数据口径经常对不上。现在的业界主流方向是Kappa架构事件源统一进消息队列流处理引擎同时承担实时计算和近实时批处理。严格说Kappa要求全量历史数据可重放如果数据量极大存储成本会明显上升。所以工业界的落地方案普遍是“批流一体”底层用同一套表格式比如Iceberg或Delta Lake用同一套引擎Flink或Spark让批处理和流处理复用同一份SQL逻辑而不是各写一套。这样一来训练集和实时特征在口径上天然对齐比在运维层面强行“对账”要省力得多。3.2 Flink落地必调的四类参数流式数据处理目前工业级最成熟的引擎是Flink。我整理了四个必调参数都是经验值按常规业务规模来给检查点间隔一般设60秒到120秒状态小的可以收到30秒太短了频繁做快照影响吞吐太长了故障恢复时间不可接受。状态后端建议用RocksDB它支持增量检查点对超大状态更友好。Kafka的Exactly-Once语义配合事务型下游写入能省掉很多“怎么又重复了”的口水仗。水位线方面处理乱序事件时要设置allowedLateness但别无脑拉大我见过有人把延迟容忍设到10分钟实时指标直接变成准实时业务方完全不能接受。3.3 实时特征拼接维表缓存的版本管理实际项目的流式特征需求往往是“实时行为离线画像”的拼接。用户实时点击、加购事件需要和离线产出的用户画像表关联工程上通常用Flink维表Join把画像数据同步到Redis或MySQL做维表缓存。这里最坑的是维表数据更新频率。如果画像每天凌晨更新一次那么维表缓存必须同步失效否则实时服务永远在用旧画像。我在项目里为这个场景专门加了“版本号”字段画像更新时递增版本号Flink作业侦测到版本变化就触发缓存刷新。实测下来比定时清缓存可靠得多至少不会出现“白天画像更新了线上还在用三天前版本”的情况。另外维表的选择也别迷信某一个组件数据量小就Redis数据量大了再考虑维度表分片这些都可以通过压测来定。4. 数据版本、元数据与特征平台模型可复现性的命根子4.1 数据湖时间旅行训练可复现的地基做AI平台和做传统数仓最大的不同就是模型复现这个硬需求。训练结果必须能回溯到它当时用的那份数据做不到这点任何实验都无法有效对比。数据湖方案里Delta Lake、Apache Iceberg、Hudi都支持时间旅行读历史快照选型时我看三个点事务能力ACID是底线不能接受小任务写到一半失败留下半成品增量读取能力流式读和增量读要高效生态兼容度团队重度用Spark就选Delta Lake舒服些多引擎混用且Flink频繁读的话Iceberg的生态更宽。数据版本不仅仅是“库表版本”还包括训练集的manifest。每次训练任务启动时训练框架要记录它消费的数据版本清单包括表名、快照时间、过滤规则随实验日志一起落库。这样才能回答“上周那个模型到底用了哪些数据”这种复盘问题。别小看这一步没有manifest你连一次有效的AB实验对比都做不出来。4.2 血缘追踪让每根数据线头清晰可见血缘记录的是数据从哪里来、经过哪些处理、被哪些任务消费。它的最大价值是排查问题时能顺着线找到根因模型指标下跌先看特征数据变了没有再看上游原始表是否异常。血缘采集常见两派解析型以解析SQL和任务定义来自动推断依赖适合Flink、Spark任务多的场景埋点型是在数据管道里显式注册父子关系适合异构脚本多的场景。我的实际做法是两者结合核心链路用埋点外围靠解析兜底。另外血缘数据本身要按天分区存储否则全量历史血缘越来越大查询性能会跟着遭殃。4.3 特征平台离线在线特征一致性的答案特征一致性是AI平台最容易翻车的地方我甚至愿意把它排在“最容易翻车榜”第一。训练时特征从离线表批量取基本都能拿到全局统计量上线时在线服务只能拿到当前时刻的局部状态。如果不做约束离线特征是“过去30天累计消费”线上实现却因为存储成本只保留最近7天行为这两个口径根本对不上模型上线不崩才怪。Feature Store要干的事不只是存特征更重要的是统一特征定义并保证离线在线的一致性。一份特征元数据离线特征存储和在线特征存储共用同一套逻辑由平台自动保证一致。另一个隐藏的关键点是时间点正确性Point-in-time Correctness生成训练集时必须避免未来数据泄漏用t1的标签预测t时刻的行为特征只能包含t时刻之前的信息。我在几个项目里都遇到过团队为了图省事直接把整月行为特征和标签一起生成等于把“答案”提前塞给了模型离线指标好看到不行线上立刻现原形。这一块是AI平台数据处理里真正体现架构师价值的地方。与其花大力气调算法参数不如先把离线在线特征一致性做扎实收益是立竿见影的。5. 数据质量监控与管道调度让数据管道自己会报警5.1 数据质量监控四类规则与阈值设置数据管道复杂到一定程度人工抽查就不现实了必须有一套自动质量监控至少覆盖四类规则。完整性表行数波动是否异常关键字段缺失率是否超阈值。及时性上游任务是否按时产出是否拖慢了训练启动。唯一性主键是否重复重复率是否在允许范围。分布漂移特征均值、分位数、类别比例是否和基线出现明显偏离。阈值设置别拍脑袋。我通常用近30天数据算出基线和标准差再设成均值加减N倍标准差的形式N一般取3。这样比手工给一个固定值更能适应业务本身的波动比如大促期间行数本来就该涨用固定阈值反而会天天误报。分布漂移一道我更推荐用近似方法算KL散度或者直方图距离精确算全量分布代价太高。5.2 调度引擎选型与依赖管理的实操经验调度层用Airflow或用云上托管调度都可以关键是几个通用原则。依赖必须显式声明上游作业完成事件触发下游而不是定时硬等定时等最容易出现“上游还没写完下游已经开跑”的脏数据问题。时区必须统一整条管道统一用UTC还是统一用东八区都行最怕调度器用UTC、业务表用本地时间凌晨窗口期到处差8小时排查起来极其痛苦。重试也不要无限重试2到3次足够重试间隔可以走指数退避。无限重试往往掩盖的是底层基础设施问题