我决定使用 Singleton 来管理此错误并刷新客户端凭据。我创建了这个类:
class ElasticsearchClientInstanceGenerator:
__instance = None
def __init__(self):
""" Virtually private constructor. """
if self.__instance is not None:
raise Exception("This class is a singleton!")
else:
ElasticsearchClientInstanceGenerator.__instance = self
self.__elasticsearch_client = None
self.__host: Optional[str] = None
self.__port: Optional[int] = None
self.__use_ssl: Optional[bool] = None
self.__verify_certs: Optional[bool] = None
def generate_elasticsearch_client(self):
import boto3
from requests_aws4auth import AWS4Auth
from elasticsearch import Elasticsearch
from elasticsearch import RequestsHttpConnection
session = boto3.Session()
credentials = session.get_credentials()
aws_auth = AWS4Auth(
credentials.access_key, credentials.secret_key, 'us-east-1', 'es', session_token=credentials.token)
self.__elasticsearch_client = Elasticsearch(
hosts=[{'host': self.__host, 'port': self.__port}],
http_auth=aws_auth, use_ssl=self.__use_ssl, verify_certs=self.__verify_certs,
connection_class=RequestsHttpConnection, timeout=30, max_retries=10, retry_on_timeout=True)
def setup(self, host: str, port: int, use_ssl: bool, verify_certs: bool):
self.__host: str = host
self.__port: int = port
self.__use_ssl: bool = use_ssl
self.__verify_certs: bool = verify_certs
@staticmethod
def get_instance():
if ElasticsearchClientInstanceGenerator.__instance is None:
ElasticsearchClientInstanceGenerator()
return ElasticsearchClientInstanceGenerator.__instance
然后我创建了这个装饰器
def try_until_succeed(func):
def catch_authorization_exception(*args, **kwargs):
for i in range(0,10):
try:
data = func(*args, **kwargs)
return data
except AuthorizationException as ae:
ElasticsearchClientInstanceGenerator.get_instance().generate_elasticsearch_client()
raise Execption('Elasticsearch general exception')
return catch_authorization_exception
我在每个需要与 Elasticcearch 连接的方法上都使用了装饰器。每当我收到 AuthorizationException 错误时,都会刷新客户端凭据,然后我可以连接到 Elasticsearch。