基于Apache Paimon与Milvus构建AI原生多模态数据湖实践 1. 项目概述当数据湖遇见AI Agent最近和几个做AI应用落地的朋友聊天大家普遍有个痛点模型本身迭代很快但喂给模型的数据管理起来却越来越像一场灾难。特别是当业务从简单的文本问答扩展到需要处理图片、音频、视频甚至传感器时序数据的多模态场景时传统的数据栈就显得力不从心了。你可能会用对象存储如S3存文件用向量数据库如Milvus存Embedding再用一个关系型数据库存元数据。数据流在多个系统间搬运、转换、对齐不仅架构复杂、延迟高更麻烦的是你很难保证一份图片文件、它的文本描述、以及由它生成的向量这三者之间的关联在任何时刻都是准确一致的。这种数据“割裂”的状态已经成为制约AI应用尤其是需要自主规划、调用工具的智能体Agent发展的主要瓶颈。这正是“AI原生多模态数据湖”要解决的问题。它不是一个新瓶装旧酒的概念而是要求数据基础设施从设计之初就将AI工作流的核心需求——特别是向量检索与多模态数据的一致性管理——作为一等公民来对待。今天我想深入聊聊的正是基于Apache Paimon和Milvus来构建这样一套基础设施的实践与思考。Paimon作为一款高性能的湖存储格式提供了流批一体、增量更新和ACID事务能力而Milvus则是业界领先的向量数据库专为海量向量检索优化。二者的结合目标直指一个核心为AI Agent构建一个统一、实时、且能理解多模态语义的数据底座。简单说就是让Agent能像我们人类一样在一个“地方”自然地关联起一段文字、一张图片和它们背后的含义并基于此做出决策。这套方案适合谁呢如果你正在或计划开发涉及复杂多模态检索的AI应用如跨模态搜索、内容推荐、智能创作、构建需要长期记忆和工具调用能力的AI Agent或者苦于现有数据平台无法支撑实时、一致的向量化数据管道那么接下来的内容或许能给你带来一些直接的参考。我们将从设计思路拆解开始一步步深入到实现细节和避坑指南。2. 核心架构设计为什么是Paimon Milvus构建AI原生数据基础设施选型是第一步也是最关键的一步。为什么是Paimon和Milvus的组合而不是其他方案这背后是对AI数据流本质需求的回应。2.1 解构AI数据流的双重需求存储与检索一个典型的面向Agent的多模态数据处理流水线可以抽象为两个核心环节统一存储层和高效检索层。这两层有截然不同的诉求试图用一个系统满足所有需求往往会导致妥协和性能瓶颈。统一存储层需要扮演“单一事实来源”的角色。它必须能容纳多模态原始数据无损存储图片、音频、视频、文本等原始文件或它们的URI以及相关的结构化元数据如创建时间、作者、标签。支持高频更新与事务AI应用的数据往往是动态的。新的数据源源不断流入旧的数据可能需要修正或删除。存储层必须支持ACID事务确保在并发写入时数据的一致性视图不被破坏。提供流批一体处理能力数据可能来自实时流如用户行为日志也可能来自批量导入如历史资料库。存储层需要能同时高效服务流式处理和批量分析任务避免维护两套系统。维护数据版本与回溯模型训练和Agent决策需要可复现性。存储层应能方便地查询数据在历史某个时间点的状态Time Travel这对于排查问题、审计溯源至关重要。高效检索层则聚焦于“智能查询”其核心是超大规模向量相似性搜索这是AI应用的基石。检索层需要能在毫秒级时间内从数十亿甚至数百亿条向量中找出与查询向量最相似的Top-K结果。支持复杂的混合查询单纯的向量搜索不够用。实际查询往往是“找到与这张图片相似且发布于上周创建者是张三的文档”。这要求检索层能同时处理向量相似度过滤和结构化属性过滤。极致的查询性能与可扩展性低延迟、高吞吐是交互式AI应用的生命线。检索层需要能通过水平扩展来应对不断增长的数据量和查询压力。2.2 技术选型逻辑各司其职与无缝衔接基于以上双重需求Paimon和Milvus的组合优势就凸显出来了。Apache Paimon作为统一的湖存储底座Paimon本质上是一个表格式Table Format类似于Apache Iceberg或Delta Lake。但它有几个特性特别契合AI场景主键表与流式更新Paimon支持定义主键并基于主键进行高效的UPSERT更新插入操作。这意味着当一份文档的元信息发生变化或生成了新的向量时你可以直接更新这条记录而无需复杂的合并操作。这对于维护数据的一致性至关重要。增量读取与流式同步Paimon的所有数据变更增、删、改都可以作为一个标准的变更数据捕获CDC流被实时读取。这为将存储层的变更实时同步到检索层提供了完美的通道。强大的生态集成Paimon与Flink深度集成可以无缝融入现有的流处理管道。同时它也可以通过Spark、Hive、StarRocks等进行查询方便了数据的批量分析与探查。Milvus作为专业的向量检索引擎Milvus是专为向量搜索而生的数据库其核心价值在于丰富的索引与量化算法支持IVF_FLAT、IVF_SQ8、HNSW等多种索引以及标量量化SQ等压缩技术能在精度和性能/成本之间提供灵活的选择。原生支持混合查询Milvus允许你在进行向量检索的同时通过布尔表达式and,or,等对标量字段即来自Paimon的元数据进行过滤一站式完成复杂查询。云原生与可扩展架构其存储计算分离、组件微服务化的架构使得扩缩容非常灵活能够轻松应对数据量和QPS的增长。组合的核心价值解耦与实时一致性这个架构最精妙之处在于“解耦”。Paimon负责可靠、一致地存储所有原始数据和元数据是数据的“源头”。Milvus则作为一个高性能的“缓存”或“索引视图”专门服务于向量检索查询。二者通过CDC流进行实时同步。这样做的好处是职责清晰每个系统做自己最擅长的事避免了单一系统的设计折衷。数据一致性有保障所有写操作都先进入Paimon这个具备事务能力的源端再异步同步到Milvus。即使Milvus出现故障或需要重建数据源始终是Paimon保证了最终的数据正确性。灵活性高你可以根据检索模式在Milvus中灵活地构建不同的向量索引例如为图片和文本分别构建索引而这些索引背后的数据都源自同一份Paimon表。注意这里有一个关键设计取舍。我们选择了“写路径统一读路径分离”。即所有写入都先到Paimon确保数据源唯一而读取时元数据查询、批量分析走Paimon低延迟的向量混合检索走Milvus。这比试图让一个系统同时承担高吞吐更新和高性能检索要现实得多。3. 构建实操从表设计到管道同步理论说清楚了我们来看具体怎么搭。假设我们要构建一个“多模态内容库”里面既有文本文档也有图片我们需要为它们生成向量并支持混合检索。3.1 数据模型与Paimon表设计首先在Paimon中设计一张主表作为所有数据的中心。这张表需要包含所有模态的元信息并为向量数据预留位置。-- 在Flink SQL中创建Paimon表 CREATE TABLE catalog.db.multimodal_assets ( asset_id STRING PRIMARY KEY NOT ENFORCED, -- 全局唯一资源ID asset_type STRING, -- 类型text, image, audio, video original_path STRING, -- 原始文件在对象存储中的路径如s3://bucket/images/001.jpg title STRING, description STRING, author STRING, tags ARRAYSTRING, -- 标签数组 created_at TIMESTAMP(3), updated_at TIMESTAMP(3), text_embedding ARRAYFLOAT, -- 文本向量可NULL image_embedding ARRAYFLOAT, -- 图像向量可NULL embedding_model STRING, -- 生成向量所用的模型名称如text-embedding-ada-002 embedding_updated_at TIMESTAMP(3) -- 向量更新时间 ) WITH ( bucket 4, -- 根据主键分桶影响并行度 bucket-key asset_id, changelog-producer full-compaction, -- 确保产生完整的CDC changelog merge-engine partial-update, -- 部分列更新适合更新向量字段 partial-update.ignore-delete true -- 忽略删除仅处理UPSERT );设计要点解析主键asset_id这是整个数据体系的锚点。无论后续是更新描述、替换文件还是更新向量都通过这个ID来定位记录。Paimon的主键表为此提供了高效的UPSERT支持。向量字段设计我们将text_embedding和image_embedding作为数组类型的列直接放在表中。这样做的好处是所有相关数据在存储层面是物理聚集的一致性由Paimon的事务保证。另一种方案是只存向量ID但那样会增加查询时的关联开销。embedding_model字段极其重要。AI模型迭代快不同版本的嵌入模型生成的向量空间不同直接比较没有意义。记录模型版本便于后续进行向量集的版本管理或重计算re-embedding。表参数partial-update当仅仅更新向量字段时这个模式可以只合并修改的列避免重写整行数据提升更新效率。3.2 构建实时向量化与同步管道数据进入Paimon表后可能是通过Flink CDC从业务库摄入或通过批量作业导入下一步是触发向量化并同步到Milvus。这里我们采用Flink作为流处理引擎来串联整个流程。// 一个简化的Flink Job示例Java API public class EmbeddingAndSyncJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(10000); // 开启检查点保证Exactly-Once语义 // 1. 从Paimon表读取变更流CDC DataStreamRowData sourceStream env.fromSource( PaimonSource.forRowData(...).table(multimodal_assets).build(), WatermarkStrategy.noWatermarks(), Paimon Source ); // 2. 处理流过滤出需要向量化的新数据或更新数据 SingleOutputStreamOperatorAssetRecord processedStream sourceStream .filter(row - needEmbedding(row)) // 自定义逻辑如asset_type更新或新增 .process(new EmbeddingProcessFunction()); // 调用Embedding API生成向量 // 3. 将生成的向量更新回Paimon表Sink processedStream.addSink( PaimonSink.forRowData(...).table(multimodal_assets).build() ); // 4. 同时将向量和关键元数据同步到Milvus processedStream.addSink( new MilvusSinkFunction() // 自定义Sink将数据写入Milvus对应Collection ); env.execute(Multimodal Embedding Sync Pipeline); } }管道核心逻辑说明变更捕获Paimon Source会持续读取multimodal_assets表的CDC日志包括INSERT和UPDATE。UPDATE可能来自对description字段的修改这可能需要重新生成文本向量。条件触发向量化在needEmbedding函数中定义触发规则。例如asset_typeimage且image_embedding为NULL的新记录或者description字段被更新且text_embedding不为NULL的旧记录需要重新生成。这里可以设计得更精细比如根据embedding_model字段判断是否需要用新模型重算。异步调用嵌入模型EmbeddingProcessFunction中需要调用外部嵌入模型API如OpenAI API、本地部署的BGE模型等。这里要注意做好错误重试、限流和降级处理避免因模型服务不稳定导致流作业失败。双写与幂等性写回Paimon将生成的向量更新到原记录的对应字段。由于是主键更新Paimon会妥善处理。写入Milvus自定义的MilvusSinkFunction需要将asset_id,text_embedding/image_embedding, 以及用于过滤的元数据如author,tags,created_at写入Milvus的Collection中。关键点在于写入Milvus的操作必须是幂等的。因为流可能会重播从检查点恢复同一条数据可能被处理多次。我们需要基于asset_id执行UPSERT操作确保Milvus中的最终状态与Paimon一致。3.3 Milvus Collection 设计与索引构建在Milvus一侧我们需要创建对应的Collection来接收数据。通常为了查询效率我们会为不同的模态或查询模式创建不同的Collection。# 使用PyMilvus创建用于文本检索的Collection from pymilvus import connections, FieldSchema, CollectionSchema, DataType, Collection, utility # 连接Milvus connections.connect(aliasdefault, hostlocalhost, port19530) # 1. 定义字段 fields [ FieldSchema(nameasset_id, dtypeDataType.VARCHAR, is_primaryTrue, max_length64), FieldSchema(nametext_embedding, dtypeDataType.FLOAT_VECTOR, dim1536), # 假设维度为1536 FieldSchema(nameauthor, dtypeDataType.VARCHAR, max_length255), FieldSchema(nametags, dtypeDataType.ARRAY, element_typeDataType.VARCHAR, max_capacity50), FieldSchema(namecreated_at, dtypeDataType.INT64), # 存储时间戳 ] schema CollectionSchema(fields, descriptionText embedding collection for multimodal assets) # 2. 创建Collection collection_name multimodal_text_assets if utility.has_collection(collection_name): utility.drop_collection(collection_name) text_collection Collection(namecollection_name, schemaschema) # 3. 创建索引 index_params { index_type: IVF_FLAT, metric_type: COSINE, # 相似度度量使用余弦相似度 params: {nlist: 1024} } text_collection.create_index(field_nametext_embedding, index_paramsindex_params) # 4. 加载Collection到内存以服务查询 text_collection.load()Milvus侧的设计考量分集合存储为文本和图片分别创建text_collection和image_collection是常见做法。这允许我们为它们配置不同的向量维度、索引参数和标量字段。虽然Milvus支持一个Collection内有多个向量字段但分开存储通常更清晰查询性能也更好优化。标量字段选择并非所有Paimon表中的元数据都需要同步到Milvus。只同步那些计划用于混合查询过滤条件的字段如author,tags,created_at。这能减少Milvus的存储和索引压力提升过滤性能。索引类型选择IVF_FLAT在精度和性能之间取得了较好平衡nlist参数需要根据数据量调整通常为sqrt(n)量级。对于十亿级别数据或对延迟极度敏感的场景可以考虑HNSW。对于存储成本敏感的场景可以使用IVF_SQ8这类量化索引。加载策略创建索引后需要将Collectionload到内存。生产环境通常使用query_node资源组和加载配置来管理多个Collection的内存占用实现按需加载或常驻内存。4. 应用层集成赋能AI Agent的查询模式基础设施搭建好后如何让上层的AI Agent方便地使用呢核心是提供一个统一的查询服务它封装了对Paimon和Milvus的调用对Agent暴露简洁的语义化接口。4.1 实现统一的多模态检索服务这个服务需要处理两类主要查询基于向量的语义检索和基于ID的详情获取。# 一个简化的检索服务示例 class MultimodalRetrievalService: def __init__(self, milvus_client, paimon_spark_session): self.text_collection milvus_client.get_collection(multimodal_text_assets) self.image_collection milvus_client.get_collection(multimodal_image_assets) self.spark paimon_spark_session def hybrid_search(self, query_vector, modalitytext, filter_exprNone, limit10): 混合检索向量相似度 标量过滤 :param query_vector: 查询向量 :param modality: 模态text 或 image :param filter_expr: Milvus布尔表达式字符串如 author 张三 and created_at 1672502400 :param limit: 返回数量 :return: 包含asset_id和分数的列表 collection self.text_collection if modality text else self.image_collection search_params {metric_type: COSINE, params: {nprobe: 20}} # nprobe影响搜索精度和速度 results collection.search( data[query_vector], anns_fieldtext_embedding if modality text else image_embedding, paramsearch_params, limitlimit, exprfilter_expr, # 这里传入混合查询的过滤条件 output_fields[asset_id] # 只返回主键ID ) # 结果格式转换 return [{id: hit.id, score: hit.score} for hit in results[0]] def get_asset_details(self, asset_ids): 根据ID列表从Paimon中获取完整的资产详情 :param asset_ids: 资产ID列表 :return: 完整的资产信息字典列表 if not asset_ids: return [] # 使用Spark SQL查询Paimon表利用主键高效点查 ids_str , .join([f{id} for id in asset_ids]) sql f SELECT asset_id, asset_type, original_path, title, description, author, tags, created_at FROM catalog.db.multimodal_assets WHERE asset_id IN ({ids_str}) df self.spark.sql(sql) return df.collect() # 返回行数据列表 def search_with_details(self, query_vector, modalitytext, filter_exprNone, limit10): 组合查询先进行向量混合检索再获取完整详情 search_results self.hybrid_search(query_vector, modality, filter_expr, limit) asset_ids [item[id] for item in search_results] details self.get_asset_details(asset_ids) # 将检索分数与详情合并 detail_map {row[asset_id]: row.asDict() for row in details} for result in search_results: result[details] detail_map.get(result[id], {}) return search_results服务设计解析两阶段查询这是性能与功能平衡的关键。第一阶段在Milvus中执行高性能的向量混合检索只返回最相关的asset_id和相似度分数。第二阶段用这些ID去Paimon中批量获取完整的元数据和原始文件路径。避免了将大量不必要的大字段如长文本、数组在向量检索时进行传输和过滤。过滤表达式filter_expr参数允许调用者传入灵活的过滤条件服务将其原样传递给Milvus。这使得Agent可以构建非常复杂的查询例如“查找与当前用户查询语义相似且标签包含‘科技’、创建于最近一个月、不是某位特定作者写的所有文章”。Spark连接Paimon这里使用Spark作为查询Paimon的引擎因为它对SQL的支持好且能方便地处理批量ID查询。对于点查Paimon的主键索引能提供高效查询。4.2 面向AI Agent的查询模式封装对于AI Agent来说它不应该关心底层是Milvus还是Paimon。我们需要提供更语义化的接口。# 面向Agent的语义化客户端 class AgenticDataClient: def __init__(self, retrieval_service, embedding_model): self.service retrieval_service self.embed_model embedding_model def search_by_text(self, query_text, modalityboth, authorNone, time_rangeNone, top_k5): Agent最常用的接口用自然语言文本进行搜索 :param query_text: 自然语言查询 :param modality: text, image, 或 both :param author: 过滤作者 :param time_range: (start_timestamp, end_timestamp) :param top_k: 返回结果数 # 1. 将查询文本向量化 query_vector self.embed_model.encode(query_text) # 2. 构建过滤表达式 filter_parts [] if author: filter_parts.append(fauthor {author}) if time_range: start_ts, end_ts time_range filter_parts.append(fcreated_at {start_ts} and created_at {end_ts}) filter_expr and .join(filter_parts) if filter_parts else None results [] if modality in [text, both]: text_results self.service.search_with_details(query_vector, text, filter_expr, top_k) results.extend(text_results) if modality in [image, both]: # 注意跨模态搜索时通常使用文本向量去搜图像向量这要求文本和图像向量在同一个对齐的空间中 image_results self.service.search_with_details(query_vector, image, filter_expr, top_k) results.extend(image_results) # 按分数排序并返回Top-K results.sort(keylambda x: x[score], reverseTrue) return results[:top_k] def get_context_for_agent(self, asset_ids, include_raw_pathFalse): 为Agent的Prompt准备上下文信息。 例如将检索到的文档的标题和描述拼接成一段文本。 details self.service.get_asset_details(asset_ids) context_parts [] for detail in details: text f标题{detail[title]}\n描述{detail[description]}\n if include_raw_path: text f原始文件{detail[original_path]}\n context_parts.append(text) return \n---\n.join(context_parts)这个AgenticDataClient对Agent非常友好。Agent只需要调用search_by_text(“寻找关于神经网络架构优化的最新图片”)就能获得结构化的、包含丰富上下文信息的结果并可以直接将这些结果注入到后续的提示词Prompt或决策逻辑中。5. 生产环境考量与避坑指南将这套方案应用到生产环境会面临许多在概念验证PoC阶段遇不到的问题。下面分享一些关键的实践经验和踩过的坑。5.1 数据一致性保障与监控“Paimon到Milvus的同步延迟”和“同步失败”是生产环境最大的风险点。必须建立完善的监控和保障机制。端到端延迟监控在数据流水线中在Paimon表写入后和Milvus写入后都打上时间戳。通过监控这两个时间戳的差值P95 P99来评估同步延迟。Flink Metrics可以很好地暴露这些指标。CDC断点续传与Exactly-Once务必开启Flink Checkpoint并确保Paimon Source和Milvus Sink都支持两阶段提交2PC或幂等写入以实现端到端的Exactly-Once语义。这意味着即使作业故障重启也不会出现数据重复或丢失。双向校验与补偿作业定期如每天运行一个离线校验作业比较Paimon中embedding_updated_at最新的N条记录是否在Milvus中存在且向量一致。如果发现不一致触发一个补偿同步作业。这里有个坑直接对比浮点数向量是否完全相等可能因为精度问题失败可以对比余弦相似度是否大于0.9999。Milvus索引重建与数据回溯当需要更换嵌入模型时所有历史向量都需要重新计算。我们的策略是在Paimon表中更新embedding_model字段标识新版本。启动一个回溯作业读取所有历史数据用新模型生成向量更新回Paimon主键更新。Paimon的CDC流会自动将更新同步到Milvus。为了不影响线上查询可以在Milvus中为新版本向量创建一个新的Collection待数据全部就绪后通过修改查询服务的配置进行切换。5.2 性能优化与成本控制随着数据量增长性能和成本问题会凸显。Paimon分区与分桶对于时间序列特征明显的数-据在Paimon表上使用PARTITIONED BY按天/月分区能极大提升按时间范围查询的效率以及过期数据清理的速度。分桶键bucket-key通常设为主键但如果是高并发更新可以考虑加入一个随机前缀来避免写热点。Milvus索引参数调优nlistIVF索引、efConstruction和MHNSW索引等参数对构建速度、查询性能和精度有巨大影响。建议在代表性数据集上进行基准测试。一个经验是随着数据量增加适当增加nlist或efConstruction以保持召回率但这会牺牲查询速度。需要在业务可接受的延迟范围内寻找平衡点。向量维度与量化评估是否可以使用维度更小的嵌入模型如从1536维降到768维这对Milvus的存储、索引内存占用和查询速度有线性级别的影响。对于精度要求稍低的场景在Milvus中使用IVF_SQ8或IVF_PQ等量化索引能用极小的精度损失换取存储和内存的大幅降低通常可压缩至原来的1/4到1/8。冷热数据分层并非所有数据都需要被高频检索。可以定义规则如“仅最近180天的数据”为热数据热数据对应的Milvus Collection常驻内存冷数据Collection则卸载release到磁盘仅在需要时加载。这需要查询服务层根据查询条件智能路由。5.3 常见问题排查实录问题向量检索结果不相关甚至乱七八糟。排查首先检查embedding_model字段。确保查询时使用的嵌入模型与库中数据生成的模型是同一个版本。不同模型生成的向量位于不同的语义空间没有可比性。这是最常见的原因。排查检查向量维度是否匹配。创建Milvus Collection时定义的dim必须与实际插入的向量维度严格一致。排查确认相似度度量标准metric_type。余弦相似度COSINE和内积IP是最常用的但需要与嵌入模型训练时使用的目标函数对齐。用错度量标准会导致排序错误。问题混合查询带过滤条件速度很慢。排查检查Milvus中用于过滤的标量字段是否创建了二级索引。对于author这类高基数字段创建Trie索引对于created_at这类范围查询字段创建STL_SORT索引可以大幅提升过滤性能。排查评估过滤条件的选择性。如果过滤后只剩很少的数据但查询时nprobe参数仍然很大意味着在大量聚类中心里搜索就会很慢。可以尝试在查询前先通过标量过滤快速缩小候选集再对这个小集合进行向量搜索即“标量过滤在前”。Milvus的expr参数执行顺序是优化的但过于复杂的表达式也可能影响性能。问题Flink同步作业消费延迟越来越大。排查检查Paimon表的写入是否产生了过多的小文件。小文件过多会导致Source读取效率低下。需要调整Paimon表的compaction相关参数如compaction.min.file-num等或者定期执行COMPACT操作来合并小文件。排查检查Embedding模型API的调用延迟和成功率。如果外部API调用缓慢或频繁失败重试会成为管道的瓶颈。需要增加并行度或引入更健壮的批处理、降级和熔断机制。问题更新了Paimon中的数据但Milvus中迟迟查不到最新状态。排查确认Paimon表的changelog-producer配置是否为full-compaction或input。none模式不会产生完整的CDC流导致更新无法被同步作业捕获。排查检查Flink作业的Checkpoint是否成功。如果Checkpoint一直失败作业可能处于不断重启的状态无法正常推进消费。排查查看Milvus Sink的日志确认写入操作是否成功是否有主键冲突等错误被忽略。构建这样一套系统是一个持续迭代和调优的过程。从Paimon和Milvus的选型与部署到数据管道和查询服务的搭建再到最终的性能优化与稳定性保障每一步都需要结合具体的业务场景和数据特性进行细致的设计。这套架构的核心优势在于其清晰的边界和强大的扩展性为应对未来更复杂的AI原生数据需求打下了坚实的基础。