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

如何使用AsyncElasticsearch建立与AWS Elasticsearch的异步连接

AWS Elasticsearch 异步连接故障解决方案

问题背景

在使用AsyncElasticsearch连接AWS Elasticsearch服务时遇到报错,同步连接逻辑可正常运行,对应同步代码如下:

from requests_aws4auth import AWS4Auth
from elasticsearch import Elasticsearch, RequestsHttpConnection

AWS_ACCESS_KEY_ID = "<id>"
AWS_SECRET_ACCESS_KEY = "<key>"
AWS_REGION = "us-east-1"
awsauth = AWS4Auth(AWS_ACCESS_KEY_ID,
                   AWS_SECRET_ACCESS_KEY,
                   AWS_REGION, 'es')


es = Elasticsearch(
        ['https://url-to-elastic-in-aws/'],
        http_auth=awsauth,
        use_ssl=True,
        verify_certs=True,
        connection_class=RequestsHttpConnection
    )

print(es.info())

注:原同步代码中错误使用await调用同步方法,已修正

异步连接代码如下:

from requests_aws4auth import AWS4Auth
from elasticsearch import AsyncElasticsearch, AIOHttpConnection

AWS_ACCESS_KEY_ID = "<id>" 
AWS_SECRET_ACCESS_KEY = "<key>" 
AWS_REGION = "us-east-1" 
awsauth = AWS4Auth(AWS_ACCESS_KEY_ID,
                   AWS_SECRET_ACCESS_KEY,
                   AWS_REGION, 'es')


es = AsyncElasticsearch(
        ['https://url-to-elastic-in-aws/'],
        http_auth=awsauth,
        use_ssl=True,
        verify_certs=True,
        connection_class=AIOHttpConnection
    )

print(await es.info())

报错分为两类:

  • 使用RequestsHttpConnection时触发TypeError: object tuple can't be used in 'await' expression
  • 使用AIOHttpConnection时触发AttributeError: 'AWS4Auth' object has no attribute 'encode'

错误原因

  • RequestsHttpConnection是同步连接类,仅适配同步requests库,方法返回值非可等待对象,不能用于异步客户端AsyncElasticsearch。
  • requests_aws4auth提供的AWS4Auth是requests专属认证类,未实现aiohttp要求的异步认证接口,因此无法被基于aiohttp的AIOHttpConnection识别。

解决方案

通过自定义兼容aiohttp的SigV4签名认证类实现异步连接,依赖botocore完成AWS签名逻辑,可运行代码如下:

前置依赖安装

pip install elasticsearch[async] botocore

完整可运行代码

import asyncio
from elasticsearch import AsyncElasticsearch, AIOHttpConnection
from botocore.auth import SigV4Auth
from botocore.awsrequest import AWSRequest
from botocore.credentials import Credentials

# 替换为自身实际配置
AWS_ACCESS_KEY_ID = "<id>"
AWS_SECRET_ACCESS_KEY = "<key>"
AWS_REGION = "us-east-1"
ES_ENDPOINT = "https://url-to-elastic-in-aws/"
SERVICE_NAME = "es"

class AWSSigV4AsyncAuth:
    def __init__(self, access_key, secret_key, region, service):
        self.credentials = Credentials(access_key, secret_key)
        self.signer = SigV4Auth(self.credentials, service, region)

    async def __call__(self, method, url, headers=None, **kwargs):
        request = AWSRequest(
            method=method,
            url=url,
            headers=headers,
            data=kwargs.get("data")
        )
        self.signer.add_auth(request)
        signed_request = request.prepare()
        headers.update(signed_request.headers)
        return (method, url, headers, kwargs)

async def main():
    auth = AWSSigV4AsyncAuth(
        AWS_ACCESS_KEY_ID,
        AWS_SECRET_ACCESS_KEY,
        AWS_REGION,
        SERVICE_NAME
    )
    es = AsyncElasticsearch(
        [ES_ENDPOINT],
        http_auth=auth,
        use_ssl=True,
        verify_certs=True,
        connection_class=AIOHttpConnection
    )
    # 测试连接
    resp = await es.info()
    print(resp)
    # 关闭客户端释放资源
    await es.close()

if __name__ == "__main__":
    asyncio.run(main())

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 00:27:00