PGVector插入新数据时的去重与数据更新技术问询
高效管理PGVector数据插入:跳过重复、更新已有ID数据
我当前使用PostgreSQL的PGVector,需要高效处理数据插入操作,核心目标是:
- 跳过已存在的数据
- 更新已有ID对应的现有数据
原代码片段
cursor.execute("SELECT EXISTS (SELECT FROM information_schema.tables WHERE table_name = 'langchain_pg_embedding')") table_exists = cursor.fetchone()[0] if not table_exists: print("Vectorstore does not exist in the database.") print("Creating Database ...") db = PGVector.from_documents( embedding=embeddings, documents=chunks, collection_name=COLLECTION_NAME, connection_string=CONNECTION_STRING ) print("Database created successfully") else: print("Vectorstore already exists in the database.") print("Checking data ...") # Check if the ID already exists in the database for chunk in chunks: cursor.execute("SELECT * FROM langchain_pg_embedding WHERE langchain_pg_embedding.cmetadata ->> 'id' = %s", (chunk.metadata["id"],)) result = cursor.fetchall() if result: print(f"ID {chunk.metadata['id']} already exists in the database.") print(result) else: print(f"Inserting ID {chunk.metadata['id']} into the database.") # Insert the chunk into the database db = PGVector.from_documents( embedding=embeddings, documents=[chunk], collection_name=COLLECTION_NAME, connection_string=CONNECTION_STRING )
原代码存在的问题
- 循环逐个检查、插入数据,批量场景下性能极差
- 仅实现了跳过重复数据,未完成更新已有ID数据的需求
- 每次插入都重新初始化PGVector实例,冗余且浪费资源
优化方案
第一步:创建唯一索引(关键前提)
先给langchain_pg_embedding表的元数据ID字段创建唯一索引,确保重复判断和更新操作的效率:
CREATE UNIQUE INDEX idx_langchain_pg_embedding_metadata_id ON langchain_pg_embedding ((cmetadata ->> 'id'));
第二步:高效批量处理代码
以下代码实现批量查询、批量插入、批量更新,同时避免冗余操作:
# 仅初始化一次PGVector实例 db = PGVector( embedding=embeddings, collection_name=COLLECTION_NAME, connection_string=CONNECTION_STRING ) cursor.execute("SELECT EXISTS (SELECT FROM information_schema.tables WHERE table_name = 'langchain_pg_embedding')") table_exists = cursor.fetchone()[0] if not table_exists: print("向量存储不存在,正在创建并插入数据...") db.add_documents(documents=chunks) print("向量存储创建完成") else: print("向量存储已存在,开始批量处理数据...") # 批量获取所有待处理数据的ID chunk_ids = [chunk.metadata["id"] for chunk in chunks] # 批量查询已存在的ID,避免循环查询 placeholders = ", ".join(["%s"] * len(chunk_ids)) cursor.execute(f"SELECT cmetadata ->> 'id' FROM langchain_pg_embedding WHERE cmetadata ->> 'id' IN ({placeholders})", chunk_ids) existing_ids = {row[0] for row in cursor.fetchall()} # 拆分待插入和待更新的数据 to_insert = [chunk for chunk in chunks if chunk.metadata["id"] not in existing_ids] to_update = [chunk for chunk in chunks if chunk.metadata["id"] in existing_ids] # 批量插入新数据 if to_insert: print(f"批量插入 {len(to_insert)} 条新数据") db.add_documents(documents=to_insert) # 批量更新已有数据 if to_update: print(f"批量更新 {len(to_update)} 条已有数据") for chunk in to_update: embedding = embeddings.embed_query(chunk.page_content) cursor.execute(""" UPDATE langchain_pg_embedding SET document = %s, embedding = %s, cmetadata = %s WHERE cmetadata ->> 'id' = %s """, (chunk.page_content, embedding, chunk.metadata, chunk.metadata["id"])) # 提交事务 cursor.connection.commit()
更简洁的UPSERT方式(一步到位)
利用PostgreSQL的ON CONFLICT语法,直接批量处理所有数据,自动跳过重复、更新已有数据:
# 仅初始化一次PGVector实例 db = PGVector( embedding=embeddings, collection_name=COLLECTION_NAME, connection_string=CONNECTION_STRING ) cursor.execute("SELECT EXISTS (SELECT FROM information_schema.tables WHERE table_name = 'langchain_pg_embedding')") table_exists = cursor.fetchone()[0] if not table_exists: print("向量存储不存在,正在创建并插入数据...") db.add_documents(documents=chunks) print("向量存储创建完成") else: print("向量存储已存在,开始批量处理数据(自动跳过/更新)...") if chunks: for chunk in chunks: embedding = embeddings.embed_query(chunk.page_content) # 用UPSERT语法实现插入或更新 cursor.execute(""" INSERT INTO langchain_pg_embedding (collection_id, document, embedding, cmetadata) VALUES ( (SELECT uuid FROM langchain_pg_collection WHERE name = %s), %s, %s, %s ) ON CONFLICT ((cmetadata ->> 'id')) DO UPDATE SET document = EXCLUDED.document, embedding = EXCLUDED.embedding, cmetadata = EXCLUDED.cmetadata """, (COLLECTION_NAME, chunk.page_content, embedding, chunk.metadata)) cursor.connection.commit() print("数据处理完成")
内容的提问来源于stack exchange,提问作者Hamza Amri
相关产品推荐
相关产品推荐

