You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PGVector插入新数据时的去重与数据更新技术问询

高效管理PGVector数据插入:跳过重复、更新已有ID数据

我当前使用PostgreSQL的PGVector,需要高效处理数据插入操作,核心目标是:

  1. 跳过已存在的数据
  2. 更新已有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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.26 15:32:42