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

如何在Python中正确使用asyncio异步获取Elasticsearch响应?

正确使用asyncio与Elasticsearch交互的方案

核心问题分析

你的代码出现“Unclosed client session”和“connection time exceed”错误,根源在于:

  • 循环内重复创建AsyncElasticsearch实例,每个实例对应独立连接会话,用完未主动关闭,导致资源泄漏。
  • 大量并发创建客户端会耗尽连接池资源,引发超时问题。

修正后的完整代码

import asyncio
from elasticsearch import AsyncElasticsearch

async def make_response(query, es):
    res = await es.search(index='index_name', query=query, size=10000)
    return res

query_list = make_queries_list_db()  # 假设该函数已正确生成查询列表

async def main():
    # 用上下文管理器创建单个Elasticsearch客户端,自动处理连接与关闭
    async with AsyncElasticsearch(port, auth) as es:
        # 批量生成协程任务,复用同一个客户端实例
        tasks = [make_response(query, es) for query in query_list]
        # 并发执行所有任务并获取结果
        list_of_res = await asyncio.gather(*tasks)
    
    # 此处可对返回的list_of_res做业务处理
    print(f"完成{len(list_of_res)}个查询请求")

asyncio.run(main())

关键细节说明

  • 客户端复用:单个AsyncElasticsearch实例内部维护连接池,能高效分配连接处理并发请求,避免重复创建连接的开销。
  • 上下文管理器:async with语法会在代码块结束时自动调用es.close(),确保所有会话和连接被正确关闭,解决“Unclosed client session”问题。
  • 超时优化:如果仍遇连接超时,可在创建客户端时调整参数:
    async with AsyncElasticsearch(
        port,
        auth,
        request_timeout=30,  # 延长请求超时时间
        max_retries=3,       # 增加重试次数
        retry_on_timeout=True  # 超时后自动重试
    ) as es:
        # 后续业务逻辑
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 16:31:16