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

OpenSearch索引脚本出现ConnectionRefusedError故障求助

问题

运行OpenSearch索引脚本时出现ConnectionRefusedError,脚本初期正常,完成约200个文件(每个含3万条文档)的索引后报错,错误栈如下:

Traceback (most recent call last):
  File "/usr/local/lib/python3.8/dist-packages/urllib3/connection.py", line 198, in _new_conn
    sock = connection.create_connection(
  File "/usr/local/lib/python3.8/dist-packages/urllib3/util/connection.py", line 85, in create_connection
    raise err
  File "/usr/local/lib/python3.8/dist-packages/urllib3/util/connection.py", line 73, in create_connection
    sock.connect(sa)
ConnectionRefusedError: [Errno 111] Connection refused

The above exception was the direct cause of the following exception:

Traceback (most recent call last):
  File "/usr/local/lib/python3.8/dist-packages/opensearchpy/connection/http_urllib3.py", line 264, in perform_request
    response = self.pool.urlopen(
  File "/usr/local/lib/python3.8/dist-packages/urllib3/connectionpool.py", line 847, in urlopen
    retries = retries.increment(
  File "/usr/local/lib/python3.8/dist-packages/urllib3/util/retry.py", line 445, in increment
    raise reraise(type(error), error, _stacktrace)
  File "/usr/local/lib/python3.8/dist-packages/urllib3/util/util.py", line 39, in reraise
    raise value
  File "/usr/local/lib/python3.8/dist-packages/urllib3/connectionpool.py", line 793, in urlopen
    response = self._make_request(
  File "/usr/local/lib/python3.8/dist-packages/urllib3/connectionpool.py", line 491, in _make_request
    raise new_e
  File "/usr/local/lib/python3.8/dist-packages/urllib3/connectionpool.py", line 467, in _make_request
    self._validate_conn(conn)
  File "/usr/local/lib/python3.8/dist-packages/urllib3/connectionpool.py", line 1099, in _validate_conn
    conn.connect()
  File "/usr/local/lib/python3.8/dist-packages/urllib3/connection.py", line 616, in connect
    self.sock = sock = self._new_conn()
  File "/usr/local/lib/python3.8/dist-packages/urllib3/connection.py", line 213, in _new_conn
    raise NewConnectionError(
urllib3.exceptions.NewConnectionError: <urllib3.connection.HTTPSConnection object at 0x7fb798541430>: Failed to establish a new connection: [Errno 111] Connection refused

During handling of the above exception, another exception occurred:

Traceback (most recent call last):
  File "indexing.py", line 162, in <module>
    create_index(index_name, path, model, batch_size=100)
  File "indexing.py", line 139, in create_index
    client.bulk(body=batch, refresh=True)
  File "/usr/local/lib/python3.8/dist-packages/opensearchpy/client/utils.py", line 181, in _wrapped
    return func(*args, params=params, headers=headers, **kwargs)
  File "/usr/local/lib/python3.8/dist-packages/opensearchpy/client/__init__.py", line 462, in bulk
    return self.transport.perform_request(
  File "/usr/local/lib/python3.8/dist-packages/opensearchpy/transport.py", line 446, in perform_request
    raise e
  File "/usr/local/lib/python3.8/dist-packages/opensearchpy/transport.py", line 409, in perform_request
    status, headers_response, data = connection.perform_request(
  File "/usr/local/lib/python3.8/dist-packages/opensearchpy/connection/http_urllib3.py", line 279, in perform_request
    raise ConnectionError("N/A", str(e), e)
opensearchpy.exceptions.ConnectionError: ConnectionError(<urllib3.connection.HTTPSConnection object at 0x7fb798541430>: Failed to establish a new connection: [Errno 111] Connection refused) caused by: NewConnectionError(<urllib3.connection.HTTPSConnection object at 0x7fb798541430>: Failed to establish a new connection: [Errno 111] Connection refused)

重启OpenSearch并从断点恢复索引后问题依旧。

环境信息

  • 操作系统:Ubuntu
  • Python版本:3.8.10
  • OpenSearch版本:2.11.1

脚本说明

脚本读取目录中的JSON文件并索引到OpenSearch集群,使用opensearchpy库交互,sentence-transformers生成嵌入向量,代码如下:

import pandas as pd
import os
import json
import torch
import time
from opensearchpy import OpenSearch
from sentence_transformers import SentenceTransformer

path = "../outputNJSONextracted"  # Directory containing your JSON files
model_card = 'sentence-transformers/msmarco-distilbert-base-tas-b'

device = torch.device("cuda" if torch.cuda.is_available() else "cpu")
print(f"Device {device}")


host = '127.0.0.1' #host = '54.93.99.186' 
port = 9200
auth = ('admin','IVIngi2024!') #('admin', 'admin')  

client = OpenSearch(
    hosts = [{'host': host, 'port': port}],
    http_auth = auth,
    use_ssl = True,
    verify_certs = False,
    ssl_assert_hostname = False,
    ssl_show_warn = False,
    timeout=30, 
    max_retries=10
)
print("Connection opened...")

index_name = 'medline-faiss-hnsw-3' # change the index name

index_body = {
 "settings": {
    "index": {
     "knn": "true",
     "refresh_interval" : -1,
     #default_pipeline": "medline-ingest-pipeline",  # embedding in script
     "number_of_shards": 5,
     "number_of_replicas": 0
    }
  },
  "mappings": {
    "properties": {
         "embedding_abstract": {
                        "type": "knn_vector",
                        "dimension": 768,
                        "method":{
                            "engine":"faiss",
                            "name": "hnsw",
                            "space_type": "innerproduct"
                        }
                       },
          "title":{"type":"text"},
          "abstract":{"type":"text"},
          "pmid":{"type":"keyword"},
          "journal":{"type":"text"},
          "pubdate":{"type":"date"},
          "authors":{"type": "text"}
    }
  }
}

response = client.indices.create(index_name, body=index_body)
print(response)

def create_index(index_name, directory_path, model, batch_size=100):
    j = 0
    documents = set()

    files_number = 0
    for filename in sorted(os.listdir(directory_path)):
        start_time = time.time()
        if filename.endswith(".json"):
            print(f"Starting indexing {filename} ...")
            # Construct the full file path
            file_path = os.path.join(directory_path, filename)
            # Read the JSON file
            with open(file_path, 'r') as file:
                # Initialize an empty list to store dictionaries
                dictionaries = []
                
                # Read the file line by line
                for line in file:
                    # Parse each line as JSON and append it to the list
                    dictionaries.append(json.loads(line))
            # Create a DataFrame
            df = pd.DataFrame(dictionaries)
            # Select only the required columns
            df = df[['pmid', 'title', 'abstract', 'journal', 'authors', 'pubdate']]
            # Output the file name
            batch = []
            for i, row in df.iterrows(): 
                pmid = row["pmid"]
                if pmid in documents:
                    continue
                else:
                    documents.add(pmid)
                    embedding = model.encode(row["abstract"])
                    doc = {
                        "pmid": pmid,
                        "abstract": row["abstract"],
                        "title": row["title"],
                        "authors": row['authors'],
                        "journal": row['journal'],
                        "pubdate": row['pubdate'],
                        "embedding_abstract": embedding
                        }
                    batch.append({"index": {"_index": index_name, "_id": pmid}})
                    batch.append(doc)
                    j += 1
                if len(batch) >= batch_size*2:
                    client.bulk(body=batch, refresh=True)
                    batch = []
                 
            if batch:
                client.bulk(body=batch, refresh=True)
                print(f"Indexed remaining documents")

            files_number += 1
            print(f"Processed file: {filename} in {time.time()-start_time}")
            print("Number of currently documents indexed ",j)
            if files_number % 100 == 0:
                print("-"*50)
                print(f"Files indexed = {files_number}")
                print()

    print("Total documents inserted = ", j)
    
    

model = SentenceTransformer(model_card)
model.to(device)
print("Creating indexing...")
start = time.time()
create_index(index_name, path, model, batch_size=100)
print(f"Time neeeded {time.time() - start}")

已采取的排查步骤

  • 将jvm.options文件中的堆内存调整至8GB(最大可用值);
  • 尝试减小批处理大小以缓解HTTP请求过大问题。

怀疑问题与大HTTP请求或OpenSearch服务器配置有关,需要Python脚本配置指导及OpenSearch索引相关的配置优化建议。


解决建议

一、Python脚本优化

  1. 添加连接重试与指数退避策略
    现有max_retries未配置重试间隔,添加指数退避避免短时间重复请求压垮服务,修改客户端初始化代码:

    from opensearchpy import OpenSearch, RequestsHttpConnection
    from requests.adapters import HTTPAdapter
    from urllib3.util.retry import Retry
    
    # 配置重试策略
    retry_strategy = Retry(
        total=10,
        backoff_factor=2,  # 指数退避,间隔为1,2,4,8...秒
        status_forcelist=[429, 500, 502, 503, 504],
        allowed_methods=["POST", "PUT"]
    )
    adapter = HTTPAdapter(max_retries=retry_strategy)
    
    client = OpenSearch(
        hosts = [{'host': host, 'port': port}],
        http_auth = auth,
        use_ssl = True,
        verify_certs = False,
        ssl_assert_hostname = False,
        ssl_show_warn = False,
        timeout=60,  # 延长超时时间
        connection_class=RequestsHttpConnection
    )
    # 挂载适配器
    client.transport.session.mount("https://", adapter)
    client.transport.session.mount("http://", adapter)
    
  2. 移除refresh=True参数
    每次bulk后强制刷新会极大增加负载,改为所有数据导入完成后手动刷新:

    # 修改bulk调用,去掉refresh=True
    client.bulk(body=batch)
    # 所有文件处理完成后执行一次刷新
    client.indices.refresh(index=index_name)
    
  3. 优化批量处理逻辑
    改用批量生成嵌入向量提升效率,增加批次大小检查:

    def create_index(index_name, directory_path, model, batch_size=100):
        j = 0
        documents = set()
        files_number = 0
        abstracts_batch = []
        metadata_batch = []
    
        for filename in sorted(os.listdir(directory_path)):
            start_time = time.time()
            if filename.endswith(".json"):
                print(f"Starting indexing {filename} ...")
                file_path = os.path.join(directory_path, filename)
                with open(file_path, 'r') as file:
                    dictionaries = [json.loads(line) for line in file]
                df = pd.DataFrame(dictionaries)
                df = df[['pmid', 'title', 'abstract', 'journal', 'authors', 'pubdate']]
                batch = []
    
                for _, row in df.iterrows():
                    pmid = row["pmid"]
                    if pmid in documents:
                        continue
                    documents.add(pmid)
                    abstracts_batch.append(row["abstract"])
                    metadata_batch.append({
                        "pmid": pmid,
                        "title": row["title"],
                        "abstract": row["abstract"],
                        "authors": row['authors'],
                        "journal": row['journal'],
                        "pubdate": row['pubdate']
                    })
                    j += 1
    
                    if len(abstracts_batch) >= batch_size:
                        embeddings = model.encode(abstracts_batch, batch_size=batch_size, device=device)
                        for meta, embedding in zip(metadata_batch, embeddings):
                            batch.append({"index": {"_index": index_name, "_id": meta["pmid"]}})
                            meta["embedding_abstract"] = embedding.tolist()
                            batch.append(meta)
                        client.bulk(body=batch)
                        abstracts_batch = []
                        metadata_batch = []
                        batch = []
    
                if abstracts_batch:
                    embeddings = model.encode(abstracts_batch, batch_size=len(abstracts_batch), device=device)
                    for meta, embedding in zip(metadata_batch, embeddings):
                        batch.append({"index": {"_index": index_name, "_id": meta["pmid"]}})
                        meta["embedding_abstract"] = embedding.tolist()
                        batch.append(meta)
                    client.bulk(body=batch)
                    print(f"Indexed remaining documents")
    
                files_number += 1
                print(f"Processed file: {filename} in {time.time()-start_time}")
                print("Number of currently documents indexed ",j)
                if files_number % 100 == 0:
                    print("-"*50)
                    print(f"Files indexed = {files_number}")
                    print()
                    time.sleep(10)  # 每处理100个文件休眠10秒
    
        client.indices.refresh(index=index_name)
        print("Total documents inserted = ", j)
    
  4. 添加异常捕获
    在bulk请求处捕获连接异常并重试:

    try:
        client.bulk(body=batch)
    except opensearchpy.exceptions.ConnectionError as e:
        print(f"Connection error occurred, retrying... {e}")
        time.sleep(5)
        client.bulk(body=batch)
    

二、OpenSearch配置优化

  1. 调整线程池设置
    修改config/opensearch.yml:

    thread_pool.bulk.queue_size: 2000
    thread_pool.bulk.size: 8  # 根据CPU核心数调整,建议为核心数的1-2倍
    
  2. 调整HTTP连接限制
    修改config/opensearch.yml:

    http.max_content_length: 100mb
    http.max_initial_line_length: 4k
    http.max_header_size: 8k
    
  3. 监控集群状态
    运行以下命令排查负载瓶颈:

    # 检查集群健康
    curl -u admin:IVIngi2024! -k https://127.0.0.1:9200/_cluster/health?pretty
    # 检查节点统计信息
    curl -u admin:IVIngi2024! -k https://127.0.0.1:9200/_nodes/stats?pretty
    # 检查线程池状态
    curl -u admin:IVIngi2024! -k https://127.0.0.1:9200/_cat/thread_pool?v
    
  4. 优化索引设置

    • 单节点环境下将number_of_shards调整为1-2;
    • 导入前关闭knn,导入完成后开启:
      # 导入前关闭knn
      curl -u admin:IVIngi2024! -k -X PUT https://127.0.0.1:9200/medline-faiss-hnsw-3/_settings -d '{"index.knn": false}' -H "Content-Type: application/json"
      # 导入完成后开启knn
      curl -u admin:IVIngi2024! -k -X PUT https://127.0.0.1:9200/medline-faiss-hnsw-3/_settings -d '{"index.knn": true}' -H "Content-Type: application/json"
      

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 07:45:06