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

Milvus批量插入向量遇Multithreaded Rendevous Error如何解决?

Milvus批量插入向量时gRPC连接异常问题

我使用以下Python脚本向Milvus集合插入约1000万条向量数据,前40批(每批10000条)运行正常,后续突然触发gRPC连接错误。

插入脚本

import numpy as np
import os
import pandas as pd
import time
from dtx_data_tools.iterate import batched, map_threaded
from pymilvus import connections, Collection

# 假设已完成客户端初始化
client = Collection("milvus_orb_benchmark")

times = []
files = [f'agg_dataset/{fp}' for fp in os.listdir('agg_dataset') if 'parquet' in fp]

db_ids_set = set()
counter = 0

for batch_file in files:
    prep_start = time.time()
    df = pd.read_parquet(batch_file).drop_duplicates(subset='scrape_uuid', keep="last")
    
    insert_ids = set(df['scrape_uuid'].tolist())

    new_uuids = insert_ids - db_ids_set
    df = df[df['scrape_uuid'].isin(new_uuids)]
    
    db_ids_set.update(new_uuids)
    
    prep_end  = time.time() - prep_start
    print(f"prep time took {prep_end} seconds")
    
    start = time.time()
    transform_time = time.time()
    df['content_vector'] = df['encoding'].apply(lambda x: np.frombuffer(x, dtype=np.float32))
    df = df[['scrape_uuid', 'content_vector']]
    print(f"Transform time: {time.time() - transform_time} seconds")
    
    try:
        client.insert(
            collection_name="milvus_orb_benchmark",
            data=df.to_dict('records')
        )
    except Exception as e:
        print(f"Insert failed: {str(e)}")
        raise
    
    counter += 10000
    print(f"Imported {counter} articles..., batch {batch_file} uploaded")
    end = time.time() - start
    print(f"batch insert_time took {end} seconds")
    times.append(end)

触发的错误

插入过程中突然抛出gRPC连接异常:

grpc._channel._InactiveRpcError: <_InactiveRpcError of RPC that terminated with:
status = StatusCode.UNAVAILABLE
details = "failed to connect to all addresses"
debug_error_string = "{"created":"@17228xxxxxx","description":"Failed to pick subchannel","file":"src/core/ext/filters/client_channel/client_channel.cc","file_line":3218,"referenced_errors":[{"created":"@17228xxxxxx","description":"failed to connect to all addresses","file":"src/core/lib/transport/error_utils.cc","file_line":165,"grpc_status":14}]}"

当前Milvus Operator配置

apiVersion: v1
kind: ServiceAccount
metadata:
  name: milvus
  annotations:
    eks.amazonaws.com/role-arn: xxxxxx
---
apiVersion: milvus.io/v1beta1
kind: Milvus
metadata:
  name: milvus
  labels:
    app: milvus
spec:
  components:
    serviceAccountName: milvus
  config:
    minio:
      bucketName: xxxxxx
      # enable AssumeRole
      useIAM: true
      useSSL: true
  dependencies:
    storage:
      external: true
      type: S3
      endpoint: xxxxxxxx
      secretRef: ""

排查与修复建议

  • 检查Milvus组件状态:执行kubectl get pods -l app=milvus,确认proxy、query、data等核心Pod是否正常运行,有无重启、CrashLoopBackoff状态。若有异常,查看Pod日志定位问题:kubectl logs <pod-name>。
  • 优化客户端gRPC配置:初始化连接时延长超时时间、增加重试次数,避免因单批请求负载过高导致连接中断:
    connections.connect(
        alias="default",
        host="your-milvus-proxy-service",
        port="19530",
        timeout=300,  # 超时时间设为5分钟
        retry_times=3
    )
    
  • 调整批量插入大小:将每批数据量从10000条降至5000条,降低单请求的资源消耗,减少连接超时概率。
  • 检查网络与资源限制:在EKS环境中,确认客户端与Milvus集群的网络带宽充足,无防火墙规则拦截gRPC端口(默认19530);检查Node节点的TCP连接数是否达到上限,可通过ss -s命令查看。
  • 调整Milvus服务端配置:在Operator配置中增加gRPC消息大小限制与超时时间,同时提升组件资源配额:
    apiVersion: milvus.io/v1beta1
    kind: Milvus
    metadata:
      name: milvus
      labels:
        app: milvus
    spec:
      components:
        serviceAccountName: milvus
        proxy:
          resources:
            requests:
              cpu: "2"
              memory: "4Gi"
            limits:
              cpu: "4"
              memory: "8Gi"
        dataNode:
          resources:
            requests:
              cpu: "4"
              memory: "8Gi"
            limits:
              cpu: "8"
              memory: "16Gi"
      config:
        minio:
          bucketName: xxxxxx
          useIAM: true
          useSSL: true
        proxy:
          grpc:
            maxRecvMsgSize: 1073741824  # 1GB
            sendMsgTimeout: 300000      # 5分钟
    
  • 添加异常重试逻辑:在插入代码块中增加重试机制,捕获gRPC异常后重试插入:
    import grpc
    
    max_retries = 3
    for attempt in range(max_retries):
        try:
            client.insert(
                collection_name="milvus_orb_benchmark",
                data=df.to_dict('records')
            )
            break
        except grpc.RpcError as e:
            if e.code() == grpc.StatusCode.UNAVAILABLE and attempt < max_retries -1:
                print(f"Connection failed, retrying attempt {attempt+1}...")
                time.sleep(5)
            else:
                raise
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 11:37:34