如何使用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
相关产品推荐
相关产品推荐

