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脚本优化
添加连接重试与指数退避策略
现有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)移除
refresh=True参数
每次bulk后强制刷新会极大增加负载,改为所有数据导入完成后手动刷新:# 修改bulk调用,去掉refresh=True client.bulk(body=batch) # 所有文件处理完成后执行一次刷新 client.indices.refresh(index=index_name)优化批量处理逻辑
改用批量生成嵌入向量提升效率,增加批次大小检查: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)添加异常捕获
在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配置优化
调整线程池设置
修改config/opensearch.yml:thread_pool.bulk.queue_size: 2000 thread_pool.bulk.size: 8 # 根据CPU核心数调整,建议为核心数的1-2倍调整HTTP连接限制
修改config/opensearch.yml:http.max_content_length: 100mb http.max_initial_line_length: 4k http.max_header_size: 8k监控集群状态
运行以下命令排查负载瓶颈:# 检查集群健康 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优化索引设置
- 单节点环境下将
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
相关产品推荐
相关产品推荐

