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
相关产品推荐
相关产品推荐

