尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

构建AI原生数据湖:Paimon与Milvus的流批一体向量化实践

构建AI原生数据湖:Paimon与Milvus的流批一体向量化实践 1. 从数据孤岛到智能体驱动为什么我们需要AI原生的数据湖如果你最近在搞AI应用尤其是涉及多模态图片、文本、音频、视频混着来或者想玩点Agent智能体的花活儿大概率会遇到一个头疼的问题数据怎么管模型训练要向量业务分析要结构化数据历史追踪要流数据这些东西散落在HDFS、对象存储、各种数据库里像个数据“百草园”。临时搭个向量库比如Milvus存一下embedding用的时候再东拼西凑流程割裂效率低下更别提让AI Agent自己去理解和使用这些分散的数据了。这背后的核心矛盾是传统大数据架构与AI原生需求之间的“代差”。传统数仓或数据湖比如基于Hive、Iceberg是为批处理和分析报表设计的核心是“存储与计算分离”和“表格式”。而AI原生应用特别是Agent需要的是低延迟、高并发的向量相似性检索以及对多模态数据的统一、实时感知与操作能力。数据不仅要存得好更要能被AI“理解”和“直接调用”。所以“AI原生多模态数据湖”这个概念就冒出来了。它不是一个简单的存储升级而是一套新的数据基础设施范式。它的目标是让多模态数据结构化、非结构化、流式在同一个地方既能高效地存储和管理作为“湖”又能被AI模型和Agent以最自然的方式尤其是通过向量实时访问、推理和操作作为“AI原生”。这就引出了标题里的两位主角Apache Paimon和Milvus。乍一看它俩好像不搭界——一个是大数据领域的流批一体存储层一个是专注向量检索的数据库。但恰恰是它们的组合指向了构建上述理想数据基础设施的一条务实路径。Paimon负责当好那个统一、高效、实时更新的“湖底”而Milvus则成为湖面上最锋利、最专业的“向量检索引擎”。两者结合不是简单拼接而是深度集成共同服务于AI应用尤其是Agentic智能体驱动的工作流。我自己的体会是现在很多团队在向量化热潮中只盯着Milvus这类向量数据库却忽略了底层数据如何实时、一致地流入向量库这个更根本的问题。结果就是向量索引陈旧Agent基于过时信息做出决策效果大打折扣。Paimon × Milvus的方案正是在尝试根治这个“数据新鲜度”的痛点让整个数据流从源头到AI应用端都是实时、统一且AI友好的。2. 核心组件拆解Paimon 与 Milvus 各自扮演什么角色要理解这个组合的威力得先抛开“112”的思维看看它们各自在AI数据栈中的独特定位以及如何互补。2.1 Apache Paimon流批一体的“湖仓底座”Paimon原名Flink Table Store是一个高性能的湖存储格式。你可以把它理解成数据湖领域的一个“新式武器”核心设计哲学是流批一体和实时更新。流批一体这意味着同一份Paimon表既可以作为流计算的源Source持续不断地吐出最新的数据变更CDC也可以作为批处理的源进行全量历史数据分析。对于AI场景这太关键了模型训练通常需要批处理全量数据而在线推理或Agent决策则需要最新的流式数据。Paimon一份存储两种消费模式避免了数据冗余和一致性难题。实时更新Paimon底层采用LSM树Log-Structured Merge-Tree结构并支持主键定义。这意味着它可以像数据库一样对记录进行高效的增、删、改操作并且这些更新能低延迟地反映到查询中。传统湖格式如Parquet是追加写更新麻烦。对于AI应用特征数据、用户画像都是动态变化的实时更新能力是刚需。统一存储层Paimon的目标是成为湖仓一体Lakehouse的存储底座。它能很好地对接Flink实时计算、Spark批处理、Hive查询等生态统一存储结构化、半结构化数据。对于多模态数据我们可以将图片、视频的元信息、路径、处理状态等结构化数据存在Paimon中而将原始文件放在对象存储如S3通过Paimon表来管理关联关系。在AI原生数据湖中的角色Paimon扮演了**统一、实时、可变更的“源数据仓库”**角色。所有业务系统的CDC数据、日志流、处理后的多模态元数据都实时流入Paimon表。它保证了数据的单一可信源并为后续的向量化处理提供了实时、干净的数据流。2.2 Milvus专为AI而生的“向量检索大脑”Milvus是一个云原生的向量数据库。它的核心能力只有一个但做到了极致海量向量的近似最近邻搜索ANN Search。为向量优化从存储将向量数据组织成便于检索的索引结构如IVF_FLAT、HNSW、计算利用GPU、SIMD指令加速向量距离计算、到分布式架构轻松横向扩展Milvus的一切设计都围绕高效向量检索展开。这是通用数据库或搜索引擎无法比拟的专业性。多向量与标量过滤除了向量Milvus也支持存储标量数据如ID、标签、时间戳。更重要的是它支持在向量相似性搜索的同时进行灵活的标量过滤比如“找到与这张图片最相似的且类别为‘猫’、上传时间在今天之后的图片”。这对于构建复杂的AI查询条件至关重要。动态数据管理Milvus支持数据的实时插入、删除和更新通过先删后插实现并能保证索引的实时可查询性取决于索引类型。这使它能够对接实时数据流为在线AI应用服务。在AI原生数据湖中的角色Milvus扮演了**AI就绪的“智能索引层”**角色。它从Paimon这样的统一数据源消费数据通过嵌入模型Embedding Model将多模态数据文本、图片等转化为向量并建立高效的向量索引。当AI应用如RAG、Agent需要根据语义进行搜索、推荐或推理时直接调用Milvus获得毫秒级的响应。2.3 组合的关键不是替代是管道化集成这里最大的误区是认为要用Milvus替代Paimon或者反之。正确的理解是管道化集成Paimon 管“全量”和“流水”所有原始和加工后的数据在这里汇聚、整理、实时更新。Milvus 管“AI视图”从Paimon的“数据流水”中实时捕捉需要向量化的数据例如新增的商品描述、用户对话日志处理后形成专门的“向量视图”。应用各取所需数据分析师直接查Paimon表做BIAI应用查Milvus做语义搜索复杂的Agent可能同时查询两者结合结构化事实和向量化语义来做决策。这种架构解耦了数据管理Paimon和AI检索Milvus让两者都能在自己擅长的领域发挥到极致并通过流式管道保持数据的同步和新鲜度。3. 构建实战如何搭建 Paimon × Milvus 的 AI 数据流水线理论说再多不如看看具体怎么搭。这里我以一个“电商多模态商品库”的场景为例拆解构建流程。目标是商品的上架、信息更新能实时反映AI客服Agent能根据用户文字描述或上传的图片实时找到最相似的商品。3.1 架构设计与数据流整体的架构数据流如下图所示概念图[业务系统] --CDC-- [Kafka] -- [Flink SQL] -- [Paimon表] (统一存储) | |--(实时流)-- [Flink Job] -- [Embedding Model] -- [Milvus] (向量索引) | |--(批查询)-- [BI工具/训练任务]数据摄入层业务数据库如MySQL的商品表变更通过CDC工具Debezium实时推送到Kafka。统一存储层Flink SQL作业消费Kafka数据写入Paimon表。这张Paimon表定义了商品的主键商品ID并包含了所有字段商品ID、名称、描述文本、图片URL、类别、价格、更新时间等。向量化处理层另一个Flink实时作业或使用Flink CDC直接监听Paimon表变更持续读取Paimon表的变更流。对于每一条新增或更新的记录它执行调用嵌入模型API如OpenAI的text-embedding-ada-002或本地部署的CLIP模型处理文本和图片将名称、描述文本和图片URL需要先下载图片转化为向量。将商品ID、向量、以及其他需要过滤的标量字段类别、价格区间组装成一条记录写入Milvus集合Collection。服务层AI应用/Agent接收用户查询文本或图片同样转化为向量向Milvus发起向量相似性搜索并可能附带标量过滤如“价格100”快速得到相关商品ID列表再根据ID回查Paimon表获取完整商品信息。数据分析直接使用Spark、Trino或Flink查询Paimon表进行传统的统计分析、报表生成。3.2 关键配置与代码片段1. 创建Paimon Catalog和表-- 在Flink SQL中创建Paimon Catalog CREATE CATALOG paimon_catalog WITH ( typepaimon, warehouses3://my-bucket/paimon/warehouse ); USE CATALOG paimon_catalog; -- 创建商品表定义主键支持CDC更新 CREATE TABLE products ( product_id BIGINT PRIMARY KEY NOT ENFORCED, name STRING, description STRING, image_url STRING, category STRING, price DECIMAL(10, 2), update_time TIMESTAMP(3) ) WITH ( bucket 4, -- 分桶数 changelog-producer full-compaction, -- 确保产生完整的changelog供下游消费 merge-engine deduplicate -- 基于主键去重实现更新 );2. 从Kafka写入Paimon-- 假设有一个Kafka topic db_products.products 包含CDC数据 CREATE TEMPORARY TABLE products_kafka (...) WITH (connectorkafka, ...); -- 使用Flink CDC的upsert模式写入Paimon自动处理INSERT/UPDATE/DELETE INSERT INTO products SELECT * FROM products_kafka;3. 构建Flink作业处理Paimon流并写入Milvus这里以Java API为例展示核心逻辑。实际生产中可能需要一个独立的Flink应用。// 1. 创建Flink表环境并连接Paimon Catalog TableEnvironment tEnv ...; tEnv.executeSql(CREATE CATALOG paimon WITH (...);); tEnv.useCatalog(paimon); // 2. 将Paimon表作为流表查询监控变更 Table productStream tEnv.sqlQuery(SELECT * FROM products /* OPTIONS(scan.modelatest) */); // 3. 转换为DataStream并应用处理函数 DataStreamRow stream tEnv.toDataStream(productStream, Row.class); stream.process(new ProcessFunctionRow, ProductVector() { private transient EmbeddingModelClient embeddingClient; Override public void open(Configuration parameters) { embeddingClient new EmbeddingModelClient(http://embedding-service:8080); } Override public void processElement(Row row, Context ctx, CollectorProductVector out) { Long productId row.getFieldAs(product_id); String name row.getFieldAs(name); String desc row.getFieldAs(description); String imageUrl row.getFieldAs(image_url); // 调用嵌入模型生成文本向量此处简化实际需处理图片 float[] textVector embeddingClient.getTextEmbedding(name desc); // 如果需要图片向量则下载图片并调用视觉模型 ProductVector pv new ProductVector(productId, textVector, row); out.collect(pv); } }) // 4. 写入Milvus Sink .addSink(new MilvusSinkFunction());4. Milvus集合Schema设计在Milvus中需要创建一个与向量对应的集合。# 使用PyMilvus from pymilvus import connections, FieldSchema, CollectionSchema, DataType, Collection, utility connections.connect(aliasdefault, hostlocalhost, port19530) fields [ FieldSchema(nameid, dtypeDataType.INT64, is_primaryTrue, auto_idFalse), FieldSchema(nameproduct_id, dtypeDataType.INT64), # 与Paimon主键对应 FieldSchema(nametext_vector, dtypeDataType.FLOAT_VECTOR, dim1536), # 向量维度 FieldSchema(namecategory, dtypeDataType.VARCHAR, max_length50), FieldSchema(nameprice, dtypeDataType.DOUBLE), ] schema CollectionSchema(fields, descriptionProduct vector collection) collection Collection(products_vector, schema) # 创建索引 index_params { index_type: IVF_FLAT, metric_type: COSINE, params: {nlist: 1024} } collection.create_index(text_vector, index_params)3.3 部署与运维要点Paimon存储规划Paimon数据存储在对象存储如S3上需要根据数据量和访问模式规划好分桶bucket策略。主键查询多的场景可以用主键做分桶字段提升点查性能。流处理作业的容错负责向量化的Flink作业需要开启Checkpoint确保从Paimon消费的位点和调用模型API的过程是容错的避免数据丢失或重复处理。Milvus集群部署生产环境建议使用Milvus集群版分离读写节点、索引节点和数据节点。根据向量规模和QPS要求进行容量规划。嵌入模型服务化将嵌入模型如Sentence Transformers, CLIP封装成高性能的gRPC或HTTP服务供Flink作业调用。需要考虑模型服务的负载均衡、批处理以提升吞吐量。监控与告警监控Paimon表的数据延迟、Flink作业的背压和消费延迟、Milvus的查询QPS/延迟/内存使用情况。设置关键指标告警。4. 面向Agentic工作流数据基础设施如何“主动”服务智能体前面的架构解决了数据“存、管、查”的问题。但对于Agentic智能体驱动应用这还不够。Agent不是被动的查询者而是主动的决策和执行者。我们的数据基础设施需要提供更高阶的能力。4.1 从“被动查询”到“主动感知与触发”在传统架构中Agent需要数据时去查Milvus或Paimon。在更高级的Agentic工作流中数据基础设施可以反过来“通知”或“触发”Agent。场景示例实时风控Agent。交易数据流写入Paimon一个实时特征工程作业基于Paimon的流数据计算用户当前交易行为的异常分数也是一个向量或标量。当分数超过阈值时这个事件本身一条高风险的记录可以实时写入一个专门的Paimon表或Kafka Topic。一个监听该数据流的Agent被自动触发获取完整的风险上下文从Paimon关联查询该用户历史行为并执行干预动作如发送验证、冻结交易。技术实现这依赖于Paimon作为流源的能力。Agent框架如LangChain, AutoGen可以集成Flink或直接读取Paimon表的CDC流将数据流转化为Agent的“观察”Observation或“事件”Event驱动其推理循环。4.2 提供“工具”Tools而非“接口”APIs对于Agent来说最好的数据接口不是一堆REST API端点而是封装好的工具Tools。我们的基础设施应该暴露“语义化”的工具。传统方式Agent调用“查询商品向量API”传入向量和过滤条件。Agentic方式我们为Agent提供一个名为search_similar_products的工具描述为“根据自然语言描述或图片寻找相似商品并可以按类别、价格筛选”。Agent在自己的工作流中可以像调用一个函数一样使用它。背后这个工具封装了1将用户输入转为向量的逻辑2调用Milvus向量检索3可能还会去Paimon补全信息。实现思路在Milvus和Paimon之上构建一层轻量的“数据服务层”可以用FastAPI等框架。这个服务层不仅提供原始数据接口更重要的是定义和实现一系列面向Agent的“工具”。这些工具通过OpenAI Function Calling、ReAct格式或LangChain Tool接口暴露给Agent框架。4.3 维护数据的“可信上下文”与“记忆”Agent在执行复杂任务时需要保持上下文Context和记忆Memory。数据湖可以成为Agent“长期记忆”的外部存储。存储交互历史将Agent与用户的完整对话历史、执行过的工具调用及其结果结构化地存入Paimon。每段对话或任务有一个会话ID作为主键。这为后续的审计、分析、以及让Agent在长对话中回顾历史提供了可能。向量化记忆检索将历史对话中的重要信息如用户偏好、已确认的事实提取出来转化为向量存入另一个Milvus集合。当Agent在新对话中遇到相关话题时可以主动从这个“记忆库”中检索相关信息实现更连贯的个性化服务。Paimon的角色提供可靠、结构化的存储保证这些记忆数据不丢失并能关联到原始业务数据如用户ID、订单ID。Milvus的角色提供从海量记忆碎片中快速进行语义检索的能力让Agent的“回忆”过程更高效。4.4 挑战与注意事项构建支持Agentic的数据基础设施挑战也随之升级数据新鲜度与一致性要求更高Agent的决策可能直接导致写操作如下单、改状态。这要求从Paimon到Milvus的数据同步延迟必须极低且需要处理好在Agent写回数据时如何同步更新两个存储的一致性例如Agent通过工具修改了商品信息这个更新需要同时写回业务库、Paimon并触发Milvus向量更新。考虑使用事务性消息或两阶段提交来保证最终一致性。工具调用的权限与审计每个暴露给Agent的工具都需要清晰的权限边界。所有工具调用必须被详细日志记录并存入Paimon用于审计。这既是安全需要也是调试和优化Agent行为的依据。成本控制实时向量化、流处理、高频的向量检索都会带来显著的计算和存储成本。需要精细化的监控对不常用的数据可以考虑冷热分离在Milvus中只保留热数据的索引。5. 性能调优与踩坑实录让流水线真正飞起来纸上谈兵容易真正跑起来坑不少。结合我自己和团队在类似项目中的经验分享几个关键的调优点和踩过的坑。5.1 Paimon层优化避免成为流水线瓶颈主键与分区键选择这是影响Paimon性能的首要因素。主键决定了更新的效率通过主键合并和点查性能。分区键通常按时间则影响数据管理过期清理和过滤查询性能。对于商品表product_id作为主键是自然的。分区键可以选择update_time的日期如dt。但注意频繁更新的数据如果分区键设计不当会导致小文件过多。踩坑记录我们曾用category作为分区键结果某些热门品类更新极频繁导致单个分区内文件数爆炸压缩合并Compact任务压力巨大影响了数据可见延迟。后来改为按天分区问题缓解。Bucket分桶数量Bucket数决定了并行度。太少的Bucket会导致单个文件过大影响读写效率太多的Bucket会产生大量小文件。建议根据数据量和集群并行度设置通常可以从“CPU核心数 * 2”开始调整。我们的经验是对于日增百万级的表Bucket数设置在16-64之间比较合适。Changelog生成模式Paimon表配置中的changelog-producer选项决定了如何为下游流作业提供变更日志。full-compaction能提供最完整的I/-U/U/-D日志但会在压缩时产生延迟。input延迟低但要求输入源本身提供完整变更日志如Kafka CDC。如果下游Flink作业需要精确的增量变更必须使用full-compaction并调整compaction.interval。小文件合并Compaction流式写入Paimon必然产生小文件。必须开启并合理配置自动压缩Compaction策略compaction.small-file-size,compaction.interval否则查询性能会随时间急剧下降。监控文件数量和平均大小是关键指标。5.2 向量化作业优化处理延迟与吞吐的平衡Flink作业并行度与背压负责调用嵌入模型API的Flink作业是典型的数据处理瓶颈。并行度设置需要与模型服务的吞吐能力匹配。过高的并行度会打垮模型服务导致超时和背压过低则数据堆积。务必监控Flink作业的背压情况。我们的做法是在模型服务前加一个缓冲队列如RedisFlink作业异步批量化地请求模型服务显著提升了吞吐。模型服务批处理绝大多数嵌入模型支持批处理输入一次处理100条文本的耗时远小于处理100次单条。将Flink流中的数据微批次如攒够50条或等待100ms后批量请求模型服务是降低延迟、提升吞吐的核心技巧。错误处理与重试网络波动、模型服务不稳定会导致偶发的嵌入失败。必须在Flink处理函数中实现健壮的重试机制如指数退避并将确实失败的消息放入死信队列Dead Letter Queue进行后续人工处理或重放避免阻塞整条流。向量维度对齐确保Paimon中不同来源的数据经过嵌入模型后产生的向量维度是一致的。例如商品标题和长描述可能用了不同的文本切片策略需要统一。否则写入Milvus时会报错。5.3 Milvus层优化检索速度与精度的艺术索引类型选择这是精度和速度的权衡。HNSW通常提供最优的查询速度和召回率但索引构建慢、内存占用高。IVF_FLAT或IVF_SQ8构建快、内存占用小查询速度也很快但需要基于数据分布进行聚类nlist参数参数调优更关键。对于亿级以上规模且内存充足的热数据HNSW是首选。对于数据量大且更新频繁的场景IVF系列更易于管理。经验分享我们一个包含2亿向量的集合使用IVF_SQ8索引nlist4096在32核机器上99%的召回率下查询延迟能稳定在10ms以内。而HNSWM16,efConstruction200能达到近100%召回率延迟在5ms内但索引大小是前者的两倍。搜索参数nprobe使用IVF索引时nprobe参数控制搜索时探查的聚类中心数。nprobe越大精度越高速度越慢。这是一个在线服务需要反复调试的参数。可以从sqrt(nlist)开始测试。标量过滤性能Milvus支持在向量搜索前后进行标量过滤。尽量使用能命中索引的标量字段进行过滤Milvus目前对标量字段的索引支持有限但合理设置仍能加速。复杂的标量过滤条件如OR逻辑、模糊匹配可能会拉低性能需要评估。内存与磁盘的平衡Milvus可以将索引和数据全部加载到内存以获得极致性能也可以部分放在磁盘。使用IVF_SQ8、SCANN等量化索引可以大幅减少内存占用。需要根据数据量、查询QPS和硬件成本做权衡。预加载与缓存对于热集合在服务启动时预加载到内存。Milvus的查询缓存缓存原始向量数据和标量结果缓存也能有效降低重复查询的延迟。5.4 端到端延迟监控与SLA保障整个流水线的SLA取决于最慢的环节。需要建立端到端的延迟监控数据新鲜度延迟从业务数据库变更到在Milvus中可被检索到之间的时间差。可以在数据中打入时间戳在Flink作业和写入Milvus后分别记录时间进行计算。查询服务延迟从AI应用发起查询到收到Milvus返回结果的时间。区分P99、P95延迟。关键链路告警对Paimon的写入延迟、Flink作业的Checkpoint失败、模型服务调用超时率、Milvus节点内存使用率等设置告警。构建这样一套系统就像维护一个精密的数据工厂。每个环节都需要精心调校监控告警必须到位。但一旦跑顺它能为上层AI应用提供稳定、实时、高质量的数据燃料那种感觉是非常畅快的。
返回列表