三亩地 三亩地SAN MU DI · CODE DIARY
ARTICLE DETAIL

日记详情

真实记录编程学习的某一天,欢迎挑你感兴趣的翻一翻。

基于Apache Paimon与Milvus构建AI原生多模态数据湖实践

基于Apache Paimon与Milvus构建AI原生多模态数据湖实践

1. 项目概述:当数据湖遇见AI Agent

最近和几个做AI应用落地的朋友聊天,大家普遍有个痛点:模型本身迭代很快,但喂给模型的数据,管理起来却越来越像一场灾难。特别是当业务从简单的文本问答,扩展到需要处理图片、音频、视频,甚至传感器时序数据的多模态场景时,传统的数据栈就显得力不从心了。你可能会用对象存储(如S3)存文件,用向量数据库(如Milvus)存Embedding,再用一个关系型数据库存元数据。数据流在多个系统间搬运、转换、对齐,不仅架构复杂、延迟高,更麻烦的是,你很难保证一份图片文件、它的文本描述、以及由它生成的向量,这三者之间的关联在任何时刻都是准确一致的。这种数据“割裂”的状态,已经成为制约AI应用,尤其是需要自主规划、调用工具的智能体(Agent)发展的主要瓶颈。

这正是“AI原生多模态数据湖”要解决的问题。它不是一个新瓶装旧酒的概念,而是要求数据基础设施从设计之初,就将AI工作流的核心需求——特别是向量检索与多模态数据的一致性管理——作为一等公民来对待。今天我想深入聊聊的,正是基于Apache PaimonMilvus来构建这样一套基础设施的实践与思考。Paimon作为一款高性能的湖存储格式,提供了流批一体、增量更新和ACID事务能力;而Milvus则是业界领先的向量数据库,专为海量向量检索优化。二者的结合,目标直指一个核心:为AI Agent构建一个统一、实时、且能理解多模态语义的数据底座。简单说,就是让Agent能像我们人类一样,在一个“地方”自然地关联起一段文字、一张图片和它们背后的含义,并基于此做出决策。

这套方案适合谁呢?如果你正在或计划开发涉及复杂多模态检索的AI应用(如跨模态搜索、内容推荐、智能创作)、构建需要长期记忆和工具调用能力的AI Agent,或者苦于现有数据平台无法支撑实时、一致的向量化数据管道,那么接下来的内容或许能给你带来一些直接的参考。我们将从设计思路拆解开始,一步步深入到实现细节和避坑指南。

2. 核心架构设计:为什么是Paimon + Milvus?

构建AI原生数据基础设施,选型是第一步,也是最关键的一步。为什么是Paimon和Milvus的组合,而不是其他方案?这背后是对AI数据流本质需求的回应。

2.1 解构AI数据流的双重需求:存储与检索

一个典型的面向Agent的多模态数据处理流水线,可以抽象为两个核心环节:统一存储层高效检索层。这两层有截然不同的诉求,试图用一个系统满足所有需求往往会导致妥协和性能瓶颈。

统一存储层需要扮演“单一事实来源”的角色。它必须能:

  1. 容纳多模态原始数据:无损存储图片、音频、视频、文本等原始文件或它们的URI,以及相关的结构化元数据(如创建时间、作者、标签)。
  2. 支持高频更新与事务:AI应用的数据往往是动态的。新的数据源源不断流入,旧的数据可能需要修正或删除。存储层必须支持ACID事务,确保在并发写入时,数据的一致性视图不被破坏。
  3. 提供流批一体处理能力:数据可能来自实时流(如用户行为日志),也可能来自批量导入(如历史资料库)。存储层需要能同时高效服务流式处理和批量分析任务,避免维护两套系统。
  4. 维护数据版本与回溯:模型训练和Agent决策需要可复现性。存储层应能方便地查询数据在历史某个时间点的状态(Time Travel),这对于排查问题、审计溯源至关重要。

高效检索层则聚焦于“智能查询”,其核心是:

  1. 超大规模向量相似性搜索:这是AI应用的基石。检索层需要能在毫秒级时间内,从数十亿甚至数百亿条向量中,找出与查询向量最相似的Top-K结果。
  2. 支持复杂的混合查询:单纯的向量搜索不够用。实际查询往往是“找到与这张图片相似,且发布于上周,创建者是张三的文档”。这要求检索层能同时处理向量相似度过滤和结构化属性过滤。
  3. 极致的查询性能与可扩展性:低延迟、高吞吐是交互式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流进行实时同步。这样做的好处是:

  1. 职责清晰:每个系统做自己最擅长的事,避免了单一系统的设计折衷。
  2. 数据一致性有保障:所有写操作都先进入Paimon这个具备事务能力的源端,再异步同步到Milvus。即使Milvus出现故障或需要重建,数据源始终是Paimon,保证了最终的数据正确性。
  3. 灵活性高:你可以根据检索模式,在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` ARRAY<STRING>, -- 标签数组 `created_at` TIMESTAMP(3), `updated_at` TIMESTAMP(3), `text_embedding` ARRAY<FLOAT>, -- 文本向量(可NULL) `image_embedding` ARRAY<FLOAT>, -- 图像向量(可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 );

设计要点解析:

  1. 主键asset_id:这是整个数据体系的锚点。无论后续是更新描述、替换文件还是更新向量,都通过这个ID来定位记录。Paimon的主键表为此提供了高效的UPSERT支持。
  2. 向量字段设计:我们将text_embeddingimage_embedding作为数组类型的列直接放在表中。这样做的好处是,所有相关数据在存储层面是物理聚集的,一致性由Paimon的事务保证。另一种方案是只存向量ID,但那样会增加查询时的关联开销。
  3. embedding_model字段极其重要。AI模型迭代快,不同版本的嵌入模型生成的向量空间不同,直接比较没有意义。记录模型版本,便于后续进行向量集的版本管理或重计算(re-embedding)。
  4. 表参数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) DataStream<RowData> sourceStream = env.fromSource( PaimonSource.forRowData(...).table("multimodal_assets").build(), WatermarkStrategy.noWatermarks(), "Paimon Source" ); // 2. 处理流:过滤出需要向量化的新数据或更新数据 SingleOutputStreamOperator<AssetRecord> 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"); } }

管道核心逻辑说明:

  1. 变更捕获:Paimon Source会持续读取multimodal_assets表的CDC日志,包括INSERTUPDATEUPDATE可能来自对description字段的修改,这可能需要重新生成文本向量。
  2. 条件触发向量化:在needEmbedding函数中定义触发规则。例如,asset_type='image'且image_embedding为NULL的新记录,或者description字段被更新且text_embedding不为NULL的旧记录(需要重新生成)。这里可以设计得更精细,比如根据embedding_model字段判断是否需要用新模型重算。
  3. 异步调用嵌入模型EmbeddingProcessFunction中需要调用外部嵌入模型API(如OpenAI API、本地部署的BGE模型等)。这里要注意做好错误重试、限流和降级处理,避免因模型服务不稳定导致流作业失败。
  4. 双写与幂等性
    • 写回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(alias="default", host='localhost', port='19530') # 1. 定义字段 fields = [ FieldSchema(name="asset_id", dtype=DataType.VARCHAR, is_primary=True, max_length=64), FieldSchema(name="text_embedding", dtype=DataType.FLOAT_VECTOR, dim=1536), # 假设维度为1536 FieldSchema(name="author", dtype=DataType.VARCHAR, max_length=255), FieldSchema(name="tags", dtype=DataType.ARRAY, element_type=DataType.VARCHAR, max_capacity=50), FieldSchema(name="created_at", dtype=DataType.INT64), # 存储时间戳 ] schema = CollectionSchema(fields, description="Text 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(name=collection_name, schema=schema) # 3. 创建索引 index_params = { "index_type": "IVF_FLAT", "metric_type": "COSINE", # 相似度度量使用余弦相似度 "params": {"nlist": 1024} } text_collection.create_index(field_name="text_embedding", index_params=index_params) # 4. 加载Collection到内存以服务查询 text_collection.load()

Milvus侧的设计考量:

  1. 分集合存储:为文本和图片分别创建text_collectionimage_collection是常见做法。这允许我们为它们配置不同的向量维度、索引参数和标量字段。虽然Milvus支持一个Collection内有多个向量字段,但分开存储通常更清晰,查询性能也更好优化。
  2. 标量字段选择:并非所有Paimon表中的元数据都需要同步到Milvus。只同步那些计划用于混合查询过滤条件的字段,如author,tags,created_at。这能减少Milvus的存储和索引压力,提升过滤性能。
  3. 索引类型选择IVF_FLAT在精度和性能之间取得了较好平衡,nlist参数需要根据数据量调整(通常为sqrt(n)量级)。对于十亿级别数据或对延迟极度敏感的场景,可以考虑HNSW。对于存储成本敏感的场景,可以使用IVF_SQ8这类量化索引。
  4. 加载策略:创建索引后,需要将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, modality='text', filter_expr=None, limit=10): """ 混合检索:向量相似度 + 标量过滤 :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_field="text_embedding" if modality == 'text' else "image_embedding", param=search_params, limit=limit, expr=filter_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, modality='text', filter_expr=None, limit=10): """ 组合查询:先进行向量混合检索,再获取完整详情 """ 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

服务设计解析:

  1. 两阶段查询:这是性能与功能平衡的关键。第一阶段,在Milvus中执行高性能的向量混合检索,只返回最相关的asset_id和相似度分数。第二阶段,用这些ID去Paimon中批量获取完整的元数据和原始文件路径。避免了将大量不必要的大字段(如长文本、数组)在向量检索时进行传输和过滤。
  2. 过滤表达式filter_expr参数允许调用者传入灵活的过滤条件,服务将其原样传递给Milvus。这使得Agent可以构建非常复杂的查询,例如“查找与当前用户查询语义相似,且标签包含‘科技’、创建于最近一个月、不是某位特定作者写的所有文章”。
  3. 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, modality='both', author=None, time_range=None, top_k=5): """ 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(f"author == '{author}'") if time_range: start_ts, end_ts = time_range filter_parts.append(f"created_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(key=lambda x: x['score'], reverse=True) return results[:top_k] def get_context_for_agent(self, asset_ids, include_raw_path=False): """ 为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的同步延迟”和“同步失败”是生产环境最大的风险点。必须建立完善的监控和保障机制。

  1. 端到端延迟监控:在数据流水线中,在Paimon表写入后和Milvus写入后都打上时间戳。通过监控这两个时间戳的差值(P95, P99)来评估同步延迟。Flink Metrics可以很好地暴露这些指标。
  2. CDC断点续传与Exactly-Once:务必开启Flink Checkpoint,并确保Paimon Source和Milvus Sink都支持两阶段提交(2PC)或幂等写入,以实现端到端的Exactly-Once语义。这意味着即使作业故障重启,也不会出现数据重复或丢失。
  3. 双向校验与补偿作业:定期(如每天)运行一个离线校验作业,比较Paimon中embedding_updated_at最新的N条记录,是否在Milvus中存在且向量一致。如果发现不一致,触发一个补偿同步作业。这里有个坑:直接对比浮点数向量是否完全相等可能因为精度问题失败,可以对比余弦相似度是否大于0.9999。
  4. Milvus索引重建与数据回溯:当需要更换嵌入模型时,所有历史向量都需要重新计算。我们的策略是:
    • 在Paimon表中更新embedding_model字段标识新版本。
    • 启动一个回溯作业,读取所有历史数据,用新模型生成向量,更新回Paimon(主键更新)。
    • Paimon的CDC流会自动将更新同步到Milvus。为了不影响线上查询,可以在Milvus中为新版本向量创建一个新的Collection,待数据全部就绪后,通过修改查询服务的配置进行切换。

5.2 性能优化与成本控制

随着数据量增长,性能和成本问题会凸显。

  1. Paimon分区与分桶:对于时间序列特征明显的数-据,在Paimon表上使用PARTITIONED BY按天/月分区,能极大提升按时间范围查询的效率以及过期数据清理的速度。分桶键(bucket-key)通常设为主键,但如果是高并发更新,可以考虑加入一个随机前缀来避免写热点。
  2. Milvus索引参数调优nlist(IVF索引)、efConstructionM(HNSW索引)等参数对构建速度、查询性能和精度有巨大影响。建议在代表性数据集上进行基准测试。一个经验是:随着数据量增加,适当增加nlistefConstruction以保持召回率,但这会牺牲查询速度。需要在业务可接受的延迟范围内寻找平衡点。
  3. 向量维度与量化:评估是否可以使用维度更小的嵌入模型(如从1536维降到768维),这对Milvus的存储、索引内存占用和查询速度有线性级别的影响。对于精度要求稍低的场景,在Milvus中使用IVF_SQ8IVF_PQ等量化索引,能用极小的精度损失换取存储和内存的大幅降低(通常可压缩至原来的1/4到1/8)。
  4. 冷热数据分层:并非所有数据都需要被高频检索。可以定义规则(如“仅最近180天的数据”为热数据),热数据对应的Milvus Collection常驻内存,冷数据Collection则卸载(release)到磁盘,仅在需要时加载。这需要查询服务层根据查询条件智能路由。

5.3 常见问题排查实录

  1. 问题:向量检索结果不相关,甚至乱七八糟。

    • 排查:首先检查embedding_model字段。确保查询时使用的嵌入模型与库中数据生成的模型是同一个版本。不同模型生成的向量位于不同的语义空间,没有可比性。这是最常见的原因。
    • 排查:检查向量维度是否匹配。创建Milvus Collection时定义的dim必须与实际插入的向量维度严格一致。
    • 排查:确认相似度度量标准(metric_type)。余弦相似度(COSINE)和内积(IP)是最常用的,但需要与嵌入模型训练时使用的目标函数对齐。用错度量标准会导致排序错误。
  2. 问题:混合查询(带过滤条件)速度很慢。

    • 排查:检查Milvus中用于过滤的标量字段是否创建了二级索引。对于author这类高基数字段,创建Trie索引;对于created_at这类范围查询字段,创建STL_SORT索引可以大幅提升过滤性能。
    • 排查:评估过滤条件的选择性。如果过滤后只剩很少的数据,但查询时nprobe参数仍然很大(意味着在大量聚类中心里搜索),就会很慢。可以尝试在查询前先通过标量过滤快速缩小候选集,再对这个小集合进行向量搜索(即“标量过滤在前”)。Milvus的expr参数执行顺序是优化的,但过于复杂的表达式也可能影响性能。
  3. 问题:Flink同步作业消费延迟越来越大。

    • 排查:检查Paimon表的写入是否产生了过多的小文件。小文件过多会导致Source读取效率低下。需要调整Paimon表的compaction相关参数(如compaction.min.file-num等),或者定期执行COMPACT操作来合并小文件。
    • 排查:检查Embedding模型API的调用延迟和成功率。如果外部API调用缓慢或频繁失败重试,会成为管道的瓶颈。需要增加并行度,或引入更健壮的批处理、降级和熔断机制。
  4. 问题:更新了Paimon中的数据,但Milvus中迟迟查不到最新状态。

    • 排查:确认Paimon表的changelog-producer配置是否为full-compactioninputnone模式不会产生完整的CDC流,导致更新无法被同步作业捕获。
    • 排查:检查Flink作业的Checkpoint是否成功。如果Checkpoint一直失败,作业可能处于不断重启的状态,无法正常推进消费。
    • 排查:查看Milvus Sink的日志,确认写入操作是否成功,是否有主键冲突等错误被忽略。

构建这样一套系统是一个持续迭代和调优的过程。从Paimon和Milvus的选型与部署,到数据管道和查询服务的搭建,再到最终的性能优化与稳定性保障,每一步都需要结合具体的业务场景和数据特性进行细致的设计。这套架构的核心优势在于其清晰的边界和强大的扩展性,为应对未来更复杂的AI原生数据需求打下了坚实的基础。

← 返回列表