1. 为什么AI让传统湖仓架构开始"不好使"了
1.1 从"存储计算分离"到"智能数据编排"的转变
先说个我自己的感受。前几年聊数据湖仓,大家挂在嘴边的还是"存储计算分离""批流一体""增量读改写"这些词。那时候的核心矛盾很简单:数据量太大,传统数仓撑不住,需要一种既能存海量明细数据、又能跑高效分析、还能兼顾实时性的架构。于是Iceberg、Hudi、Delta Lake这类开放表格式火了起来,加上S3/OSS这类廉价对象存储,配上Spark、Flink、Presto这些计算引擎,一套"湖仓一体"的架子就搭起来了。
但AI真正规模化落地之后,我发现这套架构开始出现一些"别扭"的地方。
别扭在哪?AI负载对数据的访问模式和数据仓库、传统数据分析完全不是一回事。传统BI分析是"Query-driven"——你提前知道要统计什么指标,写SQL去查,数据是结构化、模式相对固定的。AI则是"Training-driven"和"Inference-driven"——模型训练要全量扫描特征数据,而且要反复迭代,每次迭代都可能调整特征选择;模型推理要低延迟地获取特定实体的特征,甚至要实时拼接实时特征和离线特征。
这就带来几个硬约束:
- 数据形态从"表"向"特征"转变。湖里存的是明细表,但AI消费的不是明细表本身,而是经过特征工程加工后的特征宽表、特征向量、文本向量。传统湖仓的元数据管理和存储布局,对"特征"这个概念是没有原生支持的。
- 读写模式从"批扫描"向"混合模式"转变。训练任务要大批量扫描,在线推理要毫秒级点查,这俩在同一套存储上都要跑好,传统湖仓的存储层优化策略就有点捉襟见肘。
- 计算模式从"SQL-centric"向"多元计算"转变。数据科学家写Python、用Pandas/NumPy,训练框架用PyTorch/TensorFlow,特征服务要接Redis或向量数据库,这跟数仓里"一切皆SQL"的体系形成了割裂。
一句话总结:AI时代的数据湖仓,不再只是"存数据+算数据"的地方,它要变成"管理数据资产、特征资产、模型资产"的智能底座。这对我这种长期做数据平台的人来说,意味着架构思维的一个大转弯。
1.2 AI负载与传统批流负载的本质差异
为了说清楚这个差异,我拿一张表来对比:
| 对比维度 | 传统批流负载 | AI训练/推理负载 |
|---|---|---|
| 数据访问模式 | 全表扫描、按分区裁剪 | 训练全量扫描 + 推理点查/范围查 |
| 数据形态 | 结构化明细、汇总指标 | 特征向量、稀疏矩阵、文本嵌入、图结构 |
| 延迟要求 | 分钟级~小时级(批)、秒级(流) | 训练无实时要求,推理毫秒级~百毫秒级 |
| 数据更新方式 | 批量覆盖、upsert | 特征版本化、向量增量更新、模型版本回滚 |
| 计算引擎 | SQL引擎(Spark、Presto、Flink SQL) | Python生态(Pandas、PyTorch)+ 专用推理服务 |
| 元数据需求 | 表结构、分区、统计信息 | 特征血缘、特征版本、模型版本、实验记录 |
| 存储热点 | 列存格式、压缩比优先 | 混合:列存(训练)+ KV/向量索引(推理) |
这张表做完之后,我把团队拉到一起开了个会。当时我就问了一个问题:我们现在的湖仓平台,能不能做到"同一份数据,既能被Spark SQL大规模扫描做训练样本,又能被线上推理服务以毫秒级延迟查到最新特征"?
答案是不能。不是因为存储性能不够,而是整个架构从设计之初就没考虑过"特征"和"模型"这两类资产。数据只有表和文件之分,没有特征版本、模型版本的概念,更别说给它们做血缘追踪了。
所以我觉得,AI加持下的湖仓架构,本质上是一个"数据底座+特征平台+模型元数据"三层融合的架构演进。以下就是我后来落地这套架构时积累的经验和踩过的坑,希望对正在做类似改造的团队有帮助。
2. 智能湖仓架构的演进与落地形态
2.1 湖仓一体的"合体"逻辑:表格式选型是关键
先聊一个绕不开的基础选型:Iceberg、Hudi、Delta Lake三选一。
我在多个项目里都做过对比测试,也踩过不少坑。最终我个人倾向:如果团队以Spark和Flink为主、对数据湖的ACID要求高、并且未来要支撑大量并发读写的场景,Iceberg是更稳妥的选择。Hudi在数据近实时写入(比如CDC入湖)上有优势,尤其它的MOR表类型对高频小更新很友好;Delta Lake跟Spark的整合最顺滑,但如果你们的技术栈里有比较多非Spark引擎,它的开放性就是个问题。
这里我想重点说一个很多人忽略的决策维度——AI场景下的表格式选型。
AI训练读取数据时,通常希望按列裁剪、按行采样、按时间范围切片。Iceberg的hidden partitioning和manifest文件机制,在规划读取计划时能做到比较细粒度的文件裁剪,这对反复迭代的特征工程任务非常有利。Hudi的clustering功能则能帮你在数据写入后进行文件布局优化,避免大量小文件拖慢训练前的数据扫描。Delta Lake的versioned time travel能力在做训练样本回滚时很实用,你改了特征逻辑之后想回到某个历史版本的数据,直接查版本就行。
我的建议是:不要只看基准测试的跑分,要把"AI消费数据的典型路径"纳入选型矩阵里。比如我们当时有一类任务,每天要生成一份约2亿行的用户特征宽表供模型训练,同时线上推理要实时读取最近1小时的用户行为特征。我们最终是Iceberg管离线明细、Hudi管近实时增量、Redis和向量库管在线特征,三层各管一段。这个组合并不"酷",但好用。
2.2 AI时代新增的关键组件:特征存储、向量存储与模型元数据层
传统湖仓的架构图里,基本是:存储层(对象存储)、表格式层、计算引擎层、调度层、查询服务层。AI进场之后,至少还要多出三个组件,否则开发和运维会非常难受。
第一个是特征存储(Feature Store)。它解决的核心问题是:离线训练用的特征和在线推理用的特征"不一致"。这个不一致是AI项目里最隐蔽的坑——离线你用的是当天的全量特征,在线用的却是昨天的特征快照,模型上线后效果崩了但你不一定知道是特征不一致导致的。Feature Store要做的是统一特征的定义、计算逻辑和存储,离线和在线共用同一份特征逻辑,只是存储引擎不同。我们用的是Feathr和Flink实时计算特征,特征定义写在YAML里,离线和在线共用同一套定义。
第二个是向量存储(Vector Store)。大模型和RAG应用普及后,向量检索成了刚需。但向量库跟湖仓的关系,很多团队没想清楚。我的理解是:湖仓管"原始文档和结构化数据",向量库管"Embedding后的向量"。两者之间要有同步管道——当源数据在湖仓里发生变化,对应的向量也要更新。这个同步管道如果不用统一平台来调度,很快就会变成一堆没人维护的临时脚本。我们是用Dinky + Flink CDC把湖仓的增量数据实时同步到Milvus,元数据上挂血缘关系。
第三个是模型元数据层。模型本身是AI资产,但模型文件、模型版本、训练参数、评估指标、部署状态这些信息,传统湖仓完全没有管。我们搭了一个很轻量的模型元数据服务(基于MLflow二次开发),把模型注册、版本管理、线上部署状态全部接进来,模型文件本身存在对象存储上的一个独立桶里,元数据则入湖仓元数据库,让数据血缘可以一路从源表追到模型版本。
这三个组件加上去之后,湖仓架构才算真正"AI-ready"。但注意,组件越多,运维复杂度指数级上升,所以下面第5节我会专门讲踩过的一些坑。
2.3 一个可落地的分层架构示例
我自己搭过的、并且在上线后验证可行的架构,大约是下面这样。因为没用画图工具,我用文字加列表描述清楚:
存储层
- 对象存储:存原始明细数据(Iceberg表)、特征工程中间结果、模型文件
- 分布式KV/向量库:存在线特征和向量(Redis / Milvus / TiKV)
- 关系库:存平台元数据、任务实例状态、权限策略
表格式与数据管理
- Iceberg管理离线表,支持快照、分支、增量读取
- Hudi管理近实时增量流(服务于秒级到分钟级的数据新鲜度需求)
- 数据目录用独立元数据中心,统一纳管Iceberg、Hudi和向量库的元数据
计算与调度
- Spark:离线批量ETL、特征批量计算、训练样本生成
- Flink:实时特征计算、CDC增量入湖
- Doris/ClickHouse:分析型查询加速(BI和临时分析走这里)
- Airflow + DolphinScheduler:混合调度(我们最终用DolphinScheduler统一编排,因为它的告警和依赖管理做得好)
AI支撑层
- Feature Store(Feathr):特征定义、特征版本、离在线一致性管理
- 训练与推理:PyTorch + Ray(训练资源调度)+ Triton(推理服务)
- 模型元数据服务(MLflow二次开发):模型注册、版本、评估记录
- RAG支撑:文档解析管道 + Embedding管道 + Milvus向量检索
这个架构在几个典型场景下的表现是:
- 离线训练:Spark每天批量计算特征宽表->落Iceberg->Feathr注册特征->训练任务通过Feature Store API读取样本
- 在线推理:请求进来->查Redis拿用户实时特征->查Milvus拿相似向量->拼装成特征->送Triton做模型推理
- 实时特征计算:Flink消费Kafka里的行为事件->实时聚合出特征->写入Redis/Hudi->同时更新Feathr元数据
这套东西跑起来之后,最大的收益是我们的离在线特征一致性从"靠人工对口径"变成了"平台级保障"。数据团队和数据科学团队之间的"口水仗"大幅减少。但代价是平台侧的维护压力上来了,尤其是元数据服务和向量同步管道,时不时要处理一致性冲突。这个我放到踩坑部分细说。
3. AI在湖仓生命周期中的真实落地场景
3.1 元数据管理与AI自动打标
湖仓里的元数据,传统上就是表结构、分区、统计信息、权限这些。但AI来了之后,元数据管理的含义被拓宽了——我们需要给数据"语义"级别的描述,才能让模型和数据科学团队快速找到想要的数据。
我们做了两件比较实用的事。
第一件事是自动打标。以前表的标签(比如"用户行为数据""交易数据""风控相关")都是人工维护的,维护不及时,越积越乱。我们接了一个开源的大模型做语义分类,对每张表的表名、字段名、注释、前几十行样本做Embedding,然后跟已有的标签体系做语义匹配,自动生成候选标签,再由数据专员一键确认。这块覆盖率从最初的40%提升到85%以上,找表的效率明显提升。
第二件事是自动生成字段描述。很多历史表的字段没有注释,业务字段含义全靠"考古"。我们用大模型对字段名和样本值做推断,生成候选描述。准确率大概有七成,但关键是这七成不用人写,省了大量时间。剩下的三成人工修正一下就好。
这里有个实操的心得:自动打标的模型不一定要用最大参数的模型,基座模型7B~13B的量化版本就够用了,关键是Embedding质量要好。我们试过用API调用大模型和本地部署小模型两种方式,本地部署虽然初始投入高,但跑批量打标任务时不用考虑限流,成本也低得多。
3.2 智能优化与自动调优
湖仓的调度体系里,Spark作业的资源参数(executor数量、内存、并行度)以前是开发人员自己拍脑袋配的。配多了浪费资源,配少了任务跑不完甚至OOM。AI在这里可以做一件很实际的事:基于历史运行日志的智能参数推荐。
我们做了个"作业画像"系统:把每个Spark作业的历史运行数据收集起来——输入数据量、分区数、shuffle大小、执行时长、各个阶段的耗时、GC时间等,做成结构化特征,然后训练一个回归模型来预测不同资源参数下的任务执行时间,在每次任务提交前推荐一组"预计性价比最高"的参数组合。
这套系统上线之后,我们平台的整体计算资源利用率从大概30%提升到了50%多,任务失败重试的次数也降了不少。不是说AI比资深工程师更懂调优,而是人多的时候做不到每一条作业都细细调参,AI适合干这种"海量重复的精细活"。
还有一块是自动合并小文件。AI场景里小文件问题特别严重,因为特征工程经常产生各种中间结果表,一个表可能被write成几万个小文件。我们写了一个智能合并策略:根据表的文件大小分布、最近访问频率、分区新旧程度,动态决定哪些分区要触发 compaction,避开了"高峰时段压缩导致计算资源争抢"的尴尬。
3.3 数据质量与异常检测
数据质量在AI场景下有个特殊的地方:坏数据对模型效果的影响,往往不是"立即报错",而是"悄悄带偏"。比如一个特征字段因为上游bug突然大量填入默认值-1,SQL跑起来不会报错,但模型训练出来的embedding可能就扭曲了。
我们搭了一套"数据质量+特征漂移"双监控体系:
- 传统数据质量规则:非空率、唯一率、值域范围、枚举值分布等,跑在DolphinScheduler的周期任务里
- AI专项监控:特征分布的PSI(Population Stability Index)、模型输入输出的均值方差漂移检测、线上推理特征与离线训练特征的分布对比
实际操作中有一个很典型的案例:我们一个推荐模型的点击率预估特征,离线训练时用户活跃度的对数值域在0到8之间,但上线一周后线上某些新用户群的特征值大量落在8到10之间。如果只看规则校验(非空、值域),是发现不了的,因为8到10也在业务合法范围内。但监控PSI之后发现,这组特征的分布偏移超过了阈值,一查发现是新版本埋点把活跃度统计口径改了。要是没有这个监控,模型效果最少要崩两周才能被发现。
3.4 智能运维与成本治理
AI负载的另一个特点是"资源黑洞"。训练任务动辄几十上百张卡,推理服务要常驻GPU,这些资源成本比传统数仓高出好几个数量级。湖仓平台如果只管数据不管成本和资源,最后很容易变成"数据团队随便跑,财务月底一看账单崩溃"。
我们做了两个事情来治理成本:
一是资源画像与自动降配。每个业务线、每个模型、每个特征任务的资源消耗,按天做画像聚合。对连续一段时间利用率很低的资源(比如某个推理服务每秒钟只接到个位数的请求),自动触发降配提醒,或者自动把副本数缩下来。这块我们没做太智能的强化学习,就是用规则+历史分位数预测,但效果已经很可观,月成本下降了大约三成。
二是数据生命周期管理。AI训练会产生大量中间产物——临时特征表、实验数据、checkpoint文件。我们规定:临时表超过7天未访问自动进入冷归档,超过30天自动清理。checkpoint文件保留最近10个版本,旧版本自动删除。看似简单的策略,执行了半年之后,存储成本下降了大概40%。这件事难的不是写清理脚本,而是把它做成一个平台侧"无感的自动策略",让用户不需要手动申请清理,系统自动判断、提前告警、事后审计。
4. 开发规范:AI时代必须重新定义的规范清单
4.1 命名规范与分层规范
我们的数据仓库过去有一套命名规范:ods/dwd/dws/ads四层,每层有固定的前缀。AI场景进来之后,最大的问题不是层数不够,而是资产类型变多了——表、特征、模型、向量集合、实验记录,这五类资产如果混在一个命名体系里,运维会疯。
我最终定下来的命名方案层如下:
- 表仍按四层:ods(源数据)、dwd(明细)、dws(汇总)、ads(应用)
- 特征统一前缀:
feat_,特征名必须能反查特征定义文件,比如feat_user_30d_active_days - 模型统一前缀:
model_,模型名带业务域和时间,比如model_rs_click_v3_20240501 - 向量集合统一前缀:
vec_,比如vec_user_doc_embedding_v2 - 实验记录不进业务库,全部存MLflow,但与表、特征的关联关系落到元数据中心
这些规则看着简单,但推行起来阻力不小。尤其数据科学团队的同事,日常习惯用Jupyter随手跑数据,让他们遵守命名规范需要平台侧配合——比如在特征注册时强制校验前缀和格式,不合法直接拒绝注册。宁可前段卡得严一点,也比后面维护一堆"天书资产"强。
4.2 表设计规范与增量处理规范
AI场景下表设计最忌讳的一件事:训练样本表全量覆盖。
我见过太多团队每天凌晨跑一个大任务,把过去180天的训练样本重新全量算一遍,然后覆盖写入。这样做的代价是每天跑一次重活,且一旦某一天算错了,历史版本直接被冲掉,想回滚都没有退路。
我的建议是:训练样本表必须按天分区增量写入,用Iceberg/Hudi的upsert能力保留每天的快照版本。模型训练的时候通过时间旅行读取"过去某一天"的样本快照,这样即便要复现历史实验,也能拿到当时模型的训练数据全貌。
这里我贴一段简单的Iceberg时间旅行查询示例,方便大家直接抄:
-- 查询某个历史版本的表数据(基于快照ID) SELECT * FROM iceberg_catalog.db.train_sample VERSION AS OF 4789513647162953681 WHERE dt = '2024-04-01'; -- 基于时间戳回退 SELECT * FROM iceberg_catalog.db.train_sample TIMESTAMP AS OF '2024-04-01 08:00:00' WHERE dt = '2024-04-01';另外还有一条硬规定:禁止在作业里直接写"覆盖整表"的SQL。如果要重建某张表,必须走平台的重建任务,这个重建任务会自动备份旧版本再执行覆盖,防止手滑。
4.3 权限与安全模型
AI场景让权限管理变复杂了,因为数据科学家经常需要在湖仓里跑自定义Python代码,这些代码可能会接触到敏感数据。我们做了三层的控制:
第一层,数据分级。源数据统一打上密级标签,公开、内部、敏感、高敏感四个级别。不同级别的数据有不同的访问规则:高敏感数据不允许进入训练环境,敏感数据的训练任务必须走脱敏管道。
第二层,计算隔离。数据科学团队的数据探查作业,统一跑在一个受控的沙箱集群里,这个集群无法随便访问生产库的高敏感字段。需要用生产数据做特征计算时,必须走"特征申请"流程,审批通过后分配到固定表,且这个表默认脱敏屏蔽身份证号等字段。
第三层,审计追踪。所有对敏感表的访问、训练任务的数据集引用、模型的训练样本来源,全部落到审计日志里,支持回溯。这个在合规审计时特别重要。有次外部审计问我们"你们的模型X是用哪些数据训练出来的",我们花了不到半天就从数据血缘图里拉出了完整的溯源链路,这个体验比过去"翻夜摸索"好太多了。
4.4 AI资产的版本管理与评估规范
AI资产和传统数据资产一个很大的不同:它需要"版本记忆"。表被覆盖了你可以从快照找回,模型上线后效果变差了你得能快速回滚到上一版,特征的计算逻辑改了你得知道哪些历史实验是被这个改动影响的。
我定了一个规矩,称为"三件套":任何一个模型上线,必须同时在MLflow里登记三个版本:
- 模型代码版本(git commit号)
- 训练数据集版本(由特征存储的表快照ID + 数据版本组合而成)
- 特征逻辑版本(特征定义的YAML文件哈希)
这三个版本对齐之后,模型的效果问题就能分域定位:是代码变了?数据变了?还是特征逻辑变了?我们的实际排查经验是,至少一半的模型效果波动,最后都定位到"特征逻辑悄悄变了但没走变更流程"。
评估规范方面,我们要求每个模型的评估报告必须包含:离线指标(AUC、NDCG等)、线上指标(CTR、转化率等)、数据漂移指标(PSI)、特征覆盖率。且评估报告要自动归档到MLflow,不允许用"OK"这种一句话带过。这样后任模型训练同学接手时,能快速读懂前任的实验脉络。
5. 我在实际改造中踩过的坑
5.1 元数据服务的高可用问题
第一个大坑来自元数据服务。我们把表、特征、模型、向量集合的元数据全部汇到一个统一元数据中心后,这个中心就成了单点。某次上线前做故障演练,我们把元数据中心所在节点断网模拟故障,结果所有依赖它的组件全部瘫痪——作业调度起不来、特征查询失效、连模型服务的配置拉取都失败了。
后来我们做了两个调整:一是元数据中心的主库和缓存全部做双活,故障时自动切换;二是所有业务组件增加"本地降级缓存",即使元数据中心宕机,至少能在最近一次缓存版本上继续跑,不出现全平台雪崩。经验是:统一元数据是趋势,但绝对不要让它成为新的单故障点。
5.2 小文件问题在AI场景下被放大了
之前提到小文件,这里具体说说是怎么被放大的。AI特征任务经常做"按天增量写入",每跑一次就产生一批小文件;训练任务读取样本时,又按"过去180天"做跨分区扫描。如果每个分区都有大量小文件,Spark在规划stage时就容易因为文件数太多而拖慢启动阶段,甚至触发"too many open files"这类问题。
试过几种方案:设Spark自动优化参数(比如spark.sql.adaptive.enabled)、开启Iceberg的自动compaction,但效果都不够系统性。最终我们做了一个"写前规划"策略:特征任务写入前,预估分区大小,如果预估很小就"攒批"到一定体量再落盘,叫做"微批合并写入"。同时每天凌晨跑一次针对AI中间表的定向compaction。这套组合拳执行后,AI表的小文件数量减少了80%以上,特征扫描任务的执行时间缩短了约30%。
5.3 向量检索与传统SQL引擎的衔接
日子过得好好的,直到我们接RAG场景。向量数据进了Milvus之后,出现了一个尴尬问题:业务方问"能不能用SQL直接查向量?"——他们习惯用Presto或Spark SQL查数据,不愿意为了查向量再学一套API。
我们走了不少弯路,最终的做法是:把向量检索包装成Presto的自定义函数,对外提供类似vector_search(vec, k)的语法,内部实现是把请求转发给Milvus,然后拉回结果做近似度展示。这个方案一方面保留了SQL的统一入口,另一方面避免了把向量全量拉到Presto做暴力计算。当然,代价是SQL函数里不能做太复杂的过滤条件下推,经过权衡后我们接受了这个限制。
这里给一个忠告:别试图用传统数据库的索引去模拟向量检索,那种"把所有向量读出来再算余弦相似度"的做法,数据量过百万就基本没法用。除非你的场景数据量很小(几十万级),否则还是老老实实上专用向量库。
5.4 开发规范推行时的"软阻力"
技术上的坑好填,人的坑不好填。推行规范时遇到的阻力比我想象中大。数据科学团队觉得规范是"限制",数据平台团队觉得规范是"保护",这种视角差异很容易变成扯皮。
我自己试下来比较有效的方法是"平台强校验 + 默认模板 + 奖优罚劣"三管齐下:平台在提交代码和注册资产时强校验违反规范的行为,直接拒绝,这样靠个人自觉的部分大大减少;同时在自助开发页面里提供标准的特征模板、模型登记模板,让"符合规范"变成最省事的选择;最后,定期在周会上分享"规范落地带来的具体好处"——比如某次故障因为规范做得好半小时就定位了,这种正向反馈比说教有用得多。
6. 工程团队的落地路径建议与个人心得
如果你所在的团队也正准备做"AI + 湖仓"的改造,我建议不要一上来就把架构铺得很大,那是给大厂看的。务实的路径是分三步走:
第一步,先把离在线特征一致性问题解决掉。哪怕不用Feature Store,先把特征定义文件收敛到一个地方,离线和在线共用,这就能解决AI项目里最典型的一类事故。
第二步,引入向量存储但不急着大规模治理。先把RAG场景或向量检索场景跑起来,解决"怎么查向量"这个基本问题,跑通后再考虑跟湖仓的同步管道和血缘管理。
第三步,模型元数据和统一血缘做起来。这一步的价值在实验复现、审计追溯、故障排查时体现得最明显,但前期投入也最大,适合在前面两步稳定之后再启动。
我的亲身感受是,AI加持下的数据湖仓架构,真正难的不是某一项技术,而是把数据资产、特征资产、模型资产放到同一套管理体系里。这个"同一套"背后,是元数据模型的设计、是组件之间的接口约定、是开发规范的强制落地,是运维工具的自动保障。每一项单独看都不算难,合在一起面面俱到需要相当多的耐心。
最后再分享一个细节——我建议所有做数据平台的朋友,都可以给自己建一个"架构决策记录"文档,每次做技术选型、定规范、踩坑排错时,把背景、决策、教训、后续验证都写清楚。这个文档前期看起来没什么用,等团队扩到几百人、平台改了七八轮之后,它就是最宝贵的架构资产。很多坑,后人如果能从你的决策记录里提前看到,就不会再踩一遍了。