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

如何解决带交叉引用的Weaviate Collection迁移至Docker实例的报错?

Weaviate 跨实例迁移含交叉引用的 Collection 失败问题

我在将本地Weaviate实例的Collection迁移到Docker部署的Weaviate实例时,发现官方提供的迁移代码不支持交叉引用(cross-reference)属性,执行直接报错。自行尝试在查询中添加关联字段后,批量操作仍失败,但无交叉引用的Collection可正常迁移。

原代码的核心问题

  1. 硬编码了关联目标类名Object_class,无法适配实际的交叉引用目标类
  2. 直接将嵌套的关联数据塞进data_object,但Weaviate批量API要求交叉引用需单独处理,不能与普通属性混传

修改后的完整迁移代码

from typing import List, Optional
from tqdm import tqdm
from weaviate import Client

def migrate_data_from_weaviate_to_weaviate(
    client_src: Client,
    client_tgt: Client,
    from_class_name: str,
    to_class_name: str,
    from_tenant: Optional[str] = None,
    to_tenant: Optional[str] = None,
    limit: int = 1000,
    batch_size: int = 100,
    after_uuid: Optional[str] = None,
    count: int = 0,
) -> None:
    """
    跨Weaviate实例迁移数据,支持含交叉引用的Collection
    支持四种迁移模式:
            1. Class -> Class
            2. Class -> Tenant
            3. Tenant -> Class
            4. Tenant -> Tenant
    """

    # 获取源类的所有属性,分离普通属性和交叉引用属性
    source_schema = client_src.schema.get(from_class_name)
    properties = []
    cross_ref_props = []
    for prop in source_schema["properties"]:
        if prop.get("dataType") and "string" not in prop["dataType"]:
            # 识别交叉引用属性(dataType为其他类名)
            cross_ref_props.append({
                "name": prop["name"],
                "target_class": prop["dataType"][0]
            })
            # 构造查询交叉引用UUID的GraphQL片段
            properties.append(f'{prop["name"]} {{... on {prop["dataType"][0]} {{ _additional {{ id }} }} }}')
        else:
            properties.append(prop["name"])

    # 获取源类/租户的对象总数
    obj_count_query = client_src.query.aggregate(
        class_name=from_class_name
    ).with_meta_count()
    if from_tenant is not None:
        obj_count_query = obj_count_query.with_tenant(from_tenant)
    resp = obj_count_query.do()
    num_objects = resp["data"]["Aggregate"][from_class_name][0]["meta"]["count"]

    try:
        # 配置目标实例的批量API
        client_tgt.batch.configure(
            batch_size=batch_size,
            dynamic=True,
            timeout_retries=3
        )
        additional_item_config = {"tenant": to_tenant} if to_tenant else {}

        with client_tgt.batch as target_batch, tqdm(total=(num_objects - count)) as pbar:
            def ingest_data_in_batches(objects: List[dict]) -> str:
                """批量插入数据并处理交叉引用"""
                last_uuid = ""
                for obj in objects:
                    weaviate_obj = obj.copy()
                    vector = weaviate_obj["_additional"]["vector"]
                    uuid = weaviate_obj["_additional"]["id"]
                    del weaviate_obj["_additional"]

                    # 分离普通属性和交叉引用数据
                    data_object = {}
                    refs = {}
                    for prop_name, value in weaviate_obj.items():
                        is_cross_ref = any(cr["name"] == prop_name for cr in cross_ref_props)
                        if not is_cross_ref:
                            data_object[prop_name] = value
                        else:
                            # 提取关联对象的UUID
                            refs[prop_name] = [ref["_additional"]["id"] for ref in value]

                    # 插入主对象
                    if len(vector) == 0:
                        target_batch.add_data_object(
                            data_object=data_object,
                            class_name=to_class_name,
                            uuid=uuid,
                            **additional_item_config,
                        )
                    else:
                        target_batch.add_data_object(
                            data_object=data_object,
                            class_name=to_class_name,
                            uuid=uuid,
                            vector=vector,
                            **additional_item_config,
                        )

                    # 添加交叉引用
                    for prop_name, ref_uuids in refs.items():
                        target_class = next(cr["target_class"] for cr in cross_ref_props if cr["name"] == prop_name)
                        for ref_uuid in ref_uuids:
                            target_batch.add_reference(
                                from_uuid=uuid,
                                from_class_name=to_class_name,
                                from_property_name=prop_name,
                                to_uuid=ref_uuid,
                                to_class_name=target_class,
                                **additional_item_config
                            )
                    last_uuid = uuid
                return last_uuid

            # 分页迁移数据
            while True:
                query = (
                    client_src.query.get(
                        class_name=from_class_name, properties=properties
                    )
                    .with_additional(["vector", "id"])
                    .with_limit(limit)
                )
                if after_uuid:
                    query = query.with_after(after_uuid)
                if from_tenant:
                    query = query.with_tenant(from_tenant)
                source_data = query.do()

                if "errors" in source_data:
                    raise Exception(
                        f"获取数据失败,最后游标UUID: '{after_uuid}',源类: '{from_class_name}'",
                        f" 租户: '{from_tenant}'!\n" if from_tenant else "\n",
                        source_data["errors"],
                    )
                page_object = source_data["data"]["Get"][from_class_name]

                if len(page_object) == 0:
                    break
                after_uuid = ingest_data_in_batches(objects=page_object)
                pbar.update(min(limit, len(page_object)))
    except Exception as e:
        print(
            f"迁移出错,最后游标UUID: '{after_uuid}'\n"
            f"源Weaviate: 类 {from_class_name}" + (f" 租户 {from_tenant}" if from_tenant else "") + "\n"
            f"目标Weaviate: 类 {to_class_name}" + (f" 租户 {to_tenant}" if to_tenant else "") + "\n"
        )
        raise e
    finally:
        client_tgt.batch.start()

关键修改点说明

  • 动态识别交叉引用:从源类Schema自动识别交叉引用属性及目标类,无需硬编码
  • 正确查询关联UUID:查询时获取关联对象的_additional.id,这是Weaviate识别关联的唯一标识
  • 分离处理主对象与关联:先插入主对象,再用add_reference方法单独添加交叉引用,符合Weaviate批量API规范
  • 修复游标逻辑:恢复with_after参数,确保分页迁移不重复、不遗漏
  • 优化批量配置:开启动态批量和重试机制,提升迁移稳定性
  • 精准进度更新:按实际返回的对象数更新进度条,避免总数计算偏差

内容的提问来源于stack exchange,提问作者bigdwarf43

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 07:07:03