Airflow를 통해 알아보는 파이썬에서 에러난 함수를 재실행 하는 법


Airflow를 통해 알아보는 파이썬에서 에러난 함수를 재실행 하는 법

Introduction

기존 프로젝트 중 실시간 데이터의 누락을 방지 하기 위해서 파일 데이터와 비교 하며 데이터 정합성을 맞추는 배치 프로세스에서 간간히 발생하는 파일 접근 이슈로 수동으로 다시 수행 하는 경우가 생기고 있습니다.

해당 파일 데이터는 협력사가 SFTP 서버를 통해 제공 하고 있는데 네트워크 장애, 업로드 시간 지연 등 다양한 이유로 인해 파일 데이터를 얻지 못 했을 때 일어나는 상황을 어떻게 대응 해야 할지 고민 하고 있었죠.

또한, 단순히 파일 접근 이슈에만 해당 하는 것이 아닌 사내에 동작 하는 많은 배치 프로세스 모니터링, Failover 등 운영적인 측면에 대한 고민이었습니다.

그 중 첫번째로 “Apache Airflow” 도입 하여 DAG 관리와 Backfill을 이용한 Failover 전략을 고민 하게 되었고 도입 했을 때 운영 서버에서 이중화 작업, 사양에 맞는 서버 스펙 구성 등 불필요한 추가 비용 발생이나 관리 포인트 측면에서 얻는 것이 적을 것 같다는 판단이 들어 도입은 포기 하게 되었습니다.

하지만 에어플로우 내부적으로 어떻게 실패한 태스크를 재수행 하는가를 살펴본다면 충분히 Failover 전략을 구성 할 수 있을 것이라 생각을 자연스레 하게 되었는데요.


What are retry strategies in Airflow?

에어플로우는 내부적으로 메타데이터에 의해 태스크를 수행 및 관리 하고 있습니다.

에어플로우 공식 코드를 확인 해보면 아래와 같은데요.

# airflow/models/taskinstance.py
@provide_session
@internal_api_call
def _handle_failure(...):
    ...

    task_instance = _coalesce_to_orm_ti(ti=task_instance, session=session)
    failure_context = TaskInstance.fetch_handle_failure_context(
        ti=task_instance,  # type: ignore[arg-type]
        error=error,
        test_mode=test_mode,
        context=context,
        force_fail=force_fail,
        session=session,
        fail_stop=fail_stop,
    )

    _log_state(task_instance=task_instance, lead_msg="Immediate failure requested. " if force_fail else "")
    if (
        failure_context["task"]
        and failure_context["email_for_state"](failure_context["task"])
        and failure_context["task"].email
    ):
        try:
            task_instance.email_alert(error, failure_context["task"])
        except Exception:
            log.exception("Failed to send email to: %s", failure_context["task"].email)

    if failure_context["callbacks"] and failure_context["context"]:
        _run_finished_callback(
            callbacks=failure_context["callbacks"],
            context=failure_context["context"],
        )

    ...

def _run_raw_task(...):
    ...

        except (AirflowFailException, AirflowSensorTimeout) as e:
            # If AirflowFailException is raised, task should not retry.
            # If a sensor in reschedule mode reaches timeout, task should not retry.
            ti.handle_failure(e, test_mode, context, force_fail=True, session=session)  # already saves to db

def _is_eligible_to_retry(*, task_instance: TaskInstance | TaskInstancePydantic):
    """
    Is task instance is eligible for retry.

    :param task_instance: the task instance

    :meta private:
    """
    if task_instance.state == TaskInstanceState.RESTARTING:
        # If a task is cleared when running, it goes into RESTARTING state and is always
        # eligible for retry
        return True
    if not getattr(task_instance, "task", None):
        # Couldn't load the task, don't know number of retries, guess:
        return task_instance.try_number <= task_instance.max_tries

    if TYPE_CHECKING:
        assert task_instance.task

    return task_instance.task.retries and task_instance.try_number <= task_instance.max_tries

class TaskInstance(Base, LogginMixin):

  ...
  @classmethod
  def fetch_handle_failure_context(...):
  
      ...

        if force_fail or not ti.is_eligible_to_retry():
            ti.state = TaskInstanceState.FAILED
            email_for_state = operator.attrgetter("email_on_failure")
            callbacks = task.on_failure_callback if task else None

            if task and fail_stop:
                _stop_remaining_tasks(task_instance=ti, session=session)
        else:
            if ti.state == TaskInstanceState.RUNNING:
                # If the task instance is in the running state, it means it raised an exception and
                # about to retry so we record the task instance history. For other states, the task
                # instance was cleared and already recorded in the task instance history.
                from airflow.models.taskinstancehistory import TaskInstanceHistory

                TaskInstanceHistory.record_ti(ti, session=session)

            ti.state = State.UP_FOR_RETRY
            email_for_state = operator.attrgetter("email_on_retry")
            callbacks = task.on_retry_callback if task else None

        get_listener_manager().hook.on_task_instance_failed(
            previous_state=TaskInstanceState.RUNNING, task_instance=ti, error=error, session=session
        )

        return {
            "ti": ti,
            "email_for_state": email_for_state,
            "task": task,
            "callbacks": callbacks,
            "context": context,
        }

사용자가 DAG args에 입력한 retry숫자가 try_number보다 작은 경우 상태 값을 UP_FOR_RETRY로 변경 하여 재수행 하고 있었습니다.

이렇게 위 소스코드를 쭉 따라가다 보면 큰 수확을 얻지 못한 것을 알 수 있습니다.

하지만, 함수에 데코레이터로 붙어 있는 @internal_api_call 을 들어가보면 아래 처럼 tenacity 를 통해 재시도 하는 로직이 발견됩니다.

이 뿐 아닌, 에어플로우 내부 데이터베이스와 커넥션을 맺을 때도 사용 하는 코드가 보이는데요.

# airflow/api_internal/internal_api_call.py
def internal_api_call(func):
    ...

    @tenacity.retry(
        stop=tenacity.stop_after_attempt(10),
        wait=tenacity.wait_exponential(min=1),
        retry=tenacity.retry_if_exception(_is_retryable_exception),
        before_sleep=tenacity.before_log(logger, logging.WARNING),
    )
    def make_jsonrpc_request(method_name: str, params_json: str) -> bytes:
        ...

# airflow/utils/retries.py
def run_with_db_retries(max_retries: int = MAX_DB_RETRIES, logger: logging.Logger | None = None, **kwargs):
    """Return Tenacity Retrying object with project specific default."""
    import tenacity

    # Default kwargs
    retry_kwargs = dict(
        retry=tenacity.retry_if_exception_type(exception_types=(DBAPIError)),
        wait=tenacity.wait_random_exponential(multiplier=0.5, max=5),
        stop=tenacity.stop_after_attempt(max_retries),
        reraise=True,
        **kwargs,
    )
    if logger and isinstance(logger, logging.Logger):
        retry_kwargs["before_sleep"] = tenacity.before_sleep_log(logger, logging.DEBUG, True)

    return tenacity.Retrying(**retry_kwargs)

그렇다면 tenacity 를 한 번 알아봐야 되겠죠.

Understading Tenacity

general-purpose retrying library, written by Python, to simplify the task of adding retry behavior to just about anything” — 공식문서

범용적으로 사용 되는 재실행 라이브러리이며 굉장히 쉽게 구현 할 수 있다네요. 굉장히 잘 찾아온 것으로 보입니다.

다시 문제 상황으로 돌아와, 파일 데이터에 접근 하지 못 했을 때는 어떤 문제가 발생 했었을까요?

  1. 네트워크 장애로 SFTP 서버 자체에 접근 하지 못 했을 때
  2. SFTP 서버에 협력사의 파일 등록이 지연 되었을 때

등, 다양한 이유가 존재 하겠지만 가장 큰 맥락은 위 두 가지 상황입니다.

그렇다면 “특정 장애가 발생 했을 때 얼마 만큼의 시간 딜레이를 두고 몇 회 시도 하면 되겠다” 라는 결과가 나오게 됩니다.

그럼 이제 어떻게 사용 해야 하는지 공부 해 봅시다.

Install

pip3 install tenacity

Customizing Retry

Stop Conditions

  1. Fail 발생 시 몇 번을 반복 할 것인가?
from tenacity import RetryError, stop_after_attempt, retry, stop_after_delay, wait_fixed
from datetime import datetime

CALL = 0

@retry(
        stop=stop_after_attempt(1),
        )
def raise_my_exception():
    global CALL
    CALL += 1
    print(f"CALL = {CALL}, time: {datetime.now()}")
    raise ValueError(f"{CALL} has Value Error")

try:
    raise_my_exception()
except RetryError as e:
    print(f"retry error = {e}")

stats = raise_my_exception.retry.statistics
print(stats)
print(f"Retry 시작 시간: {datetime.now()}")
print(f"시도 횟수: {stats['attempt_number']}")

"""
output

CALL = 1, time: 2024-10-17 14:30:17.043717
CALL = 2, time: 2024-10-17 14:30:17.043843
CALL = 3, time: 2024-10-17 14:30:17.043870
retry error = RetryError[<Future at 0x7fe5a00aa860 state=finished raised ValueError>]
{'start_time': 103128.486361125, 'attempt_number': 3, 'idle_for': 0, 'delay_since_first_attempt': 0.00016700000560376793}
Retry 시작 시간: 2024-10-17 14:30:17.043921
시도 횟수: 3
""

2. Fail 발생 시 몇 초간 반복 할 것인가?

from tenacity import RetryError, stop_after_attempt, retry, stop_after_delay, wait_fixed
from datetime import datetime

CALL = 0

@retry(
        stop=stop_after_delay(3),
        wait=wait_fixed(1),  # NOTE: 3초간 무한 반복 요청을 방지 하기 위해 sleep 추가
        )
def raise_my_exception():
    global CALL
    CALL += 1
    print(f"CALL = {CALL}, time: {datetime.now()}")
    raise ValueError(f"{CALL} has Value Error")

try:
    raise_my_exception()
except RetryError as e:
    print(f"retry error = {e}")

stats = raise_my_exception.retry.statistics

print(f"Retry 시작 시간: {datetime.now()}")
print(f"시도 횟수: {stats['attempt_number']}")
print(f"함수 첫번째 실행 후 총 경과 시간: {stats['delay_since_first_attempt']}")

""" output

CALL = 1, time: 2024-10-17 14:36:42.601391
CALL = 2, time: 2024-10-17 14:36:43.606668
CALL = 3, time: 2024-10-17 14:36:44.612161
CALL = 4, time: 2024-10-17 14:36:45.614032
retry error = RetryError[<Future at 0x7faaa8362ce0 state=finished raised ValueError>]
Retry 시작 시간: 2024-10-17 14:36:45.614674
시도 횟수: 4
함수 첫번째 실행 후 총 경과 시간: 3.0128697080072016
"""

3. Fail발생 시 몇초 동안 몇 회를 반복 할 것인가? (1, 2 조합)

from tenacity import RetryError, stop_after_attempt, retry, stop_after_delay, wait_fixed
from datetime import datetime

CALL = 0

@retry(
        stop=(stop_after_delay(3) | stop_after_attempt(5)),
        wait=wait_fixed(1),
        )
def raise_my_exception():
    global CALL
    CALL += 1
    print(f"CALL = {CALL}, time: {datetime.now()}")
    raise ValueError(f"{CALL} has Value Error")

try:
    raise_my_exception()
except RetryError as e:
    print(f"retry error = {e}")

stats = raise_my_exception.retry.statistics

print(f"Retry 시작 시간: {datetime.now()}")
print(f"시도 횟수: {stats['attempt_number']}")
print(f"함수 첫번째 실행 후 총 경과 시간: {stats['delay_since_first_attempt']}")

""" output

CALL = 1, time: 2024-10-17 14:41:27.448525
CALL = 2, time: 2024-10-17 14:41:28.452302
CALL = 3, time: 2024-10-17 14:41:29.455972
CALL = 4, time: 2024-10-17 14:41:30.461291
CALL = 5, time: 2024-10-17 14:41:31.466631
retry error = RetryError[<Future at 0x7faa90778190 state=finished raised ValueError>]
Retry 시작 시간: 2024-10-17 14:41:31.467209
시도 횟수: 5
함수 첫번째 실행 후 총 경과 시간: 4.018330582999624
"""

3번 조합을 사용 할 땐 둘의 조건이 OR로 이루어져있다는 것을 명심 해야 합니다. 둘 중 하나의 조건이 만족한다면 끝나기 때문이죠.

Retry Strategy

  1. Fail 발생 시 특정 Exception일 경우 Retry
from tenacity import RetryError, stop_after_attempt, retry, retry_if_exception
from datetime import datetime

CALL = 0

class ExampleException(Exception):
    pass

@retry(
        retry=retry_if_exception(ExampleException),
        stop=stop_after_attempt(3),
        )
def raise_my_exception():
    global CALL
    CALL += 1
    print(f"CALL = {CALL}, time: {datetime.now()}")
    raise ExampleException(f"{CALL} has Value Error")

try:
    raise_my_exception()
except RetryError as e:
    print(f"retry error = {e}")

stats = raise_my_exception.retry.statistics

print(f"Retry 시작 시간: {datetime.now()}")
print(f"시도 횟수: {stats['attempt_number']}")
print(f"함수 첫번째 실행 후 총 경과 시간: {stats['delay_since_first_attempt']}")

""" output
CALL = 1, time: 2024-10-17 15:15:30.594013
CALL = 2, time: 2024-10-17 15:15:30.594111
CALL = 3, time: 2024-10-17 15:15:30.594137
retry error = RetryError[<Future at 0x7fc85845b040 state=finished raised ExampleException>]
Retry 시작 시간: 2024-10-17 15:15:30.594166
시도 횟수: 3
함수 첫번째 실행 후 총 경과 시간: 0.00013662499259226024
"""

Exception 이 달라도 RetryError 로 Catch가 되는데 특정 Exception 만 잡고 싶다면 reraise=True 옵션을 추가로 적용 해야합니다.

**reraise**옵션 추가

from tenacity import RetryError, stop_after_attempt, retry, retry_if_exception
from datetime import datetime

CALL = 0

class ExampleException(Exception):
    pass

@retry(
        retry=retry_if_exception(ExampleException),
        stop=stop_after_attempt(3),
        reraise=True,
        )
def raise_my_exception():
    global CALL
    CALL += 1
    print(f"CALL = {CALL}, time: {datetime.now()}")
    raise ExampleException(f"{CALL} has Example Exception Error")

try:
    raise_my_exception()
except RetryError as e:
    print(f"retry error = {e}")
except ExampleException as e:
    print(f"example exception error = {e}")

stats = raise_my_exception.retry.statistics

print(f"Retry 시작 시간: {datetime.now()}")
print(f"시도 횟수: {stats['attempt_number']}")
print(f"함수 첫번째 실행 후 총 경과 시간: {stats['delay_since_first_attempt']}")

"""
CALL = 1, time: 2024-10-17 15:27:27.223505
CALL = 2, time: 2024-10-17 15:27:27.223609
CALL = 3, time: 2024-10-17 15:27:27.223635
example exception error = 3 has Example Exception Error
Retry 시작 시간: 2024-10-17 15:27:27.223662
시도 횟수: 3
함수 첫번째 실행 후 총 경과 시간: 0.00014320900663733482
"""

2. 해당 함수의 결과 값이 None인 경우 Retry

from tenacity import retry, retry_if_result

def is_none(value):
    return value is None

@retry(retry=retry_if_result(is_none), stop=stop_after_attempt(3))
def return_none():
    print("return None")

return_none()

"""
return None
return None
return None
Traceback (most recent call last):
...
tenacity.RetryError: RetryError[<Future at 0x7f9928578850 state=finished returned NoneType>]
"""

retry 데코레이터 래핑 영역에서 수행 할 함수의 결과 값이 retry_if_result 에 들어갈 일급 함수의 반환 값과 같다면(filter 처럼 Conditional value를 반환 하여 사용 하는게 현명할 것 같네요) retry프로세스가 수행 됩니다.

Logging

  1. Before Sleep Logging
import logging
import sys

from tenacity import RetryError, stop_after_attempt, retry, before_sleep_log

logging.basicConfig(stream=sys.stderr, level=logging.DEBUG)

logger = logging.getLogger(__name__)

CALL = 0

@retry(
        stop=stop_after_attempt(3),
        before_sleep=before_sleep_log(logger, logging.DEBUG),
        )
def raise_my_exception():
    global CALL
    CALL += 1
    print(f"CALL = {CALL}, time: {datetime.now()}")
    raise ValueError(f"{CALL} has Value Error")

try:
    raise_my_exception()
except RetryError as e:
    print(f"retry error = {e}")

stats = raise_my_exception.retry.statistics

print(f"Retry 시작 시간: {datetime.now()}")
print(f"시도 횟수: {stats['attempt_number']}")
print(f"함수 첫번째 실행 후 총 경과 시간: {stats['delay_since_first_attempt']}")

""" output
CALL = 1, time: 2024-10-17 15:45:47.315753
DEBUG:__main__:Retrying __main__.raise_my_exception in 0.0 seconds as it raised ValueError: 1 has Value Error.
CALL = 2, time: 2024-10-17 15:45:47.315999
DEBUG:__main__:Retrying __main__.raise_my_exception in 0.0 seconds as it raised ValueError: 2 has Value Error.
CALL = 3, time: 2024-10-17 15:45:47.316059
retry error = RetryError[<Future at 0x7f8ea037f970 state=finished raised ValueError>]
Retry 시작 시간: 2024-10-17 15:45:47.316145
시도 횟수: 3
함수 첫번째 실행 후 총 경과 시간: 0.00036025000736117363
"""

2. After Logging

import logging
import sys

from tenacity import RetryError, stop_after_attempt, retry, after_log, before_sleep_log

logging.basicConfig(stream=sys.stderr, level=logging.DEBUG)

logger = logging.getLogger(__name__)

CALL = 0

@retry(
        stop=stop_after_attempt(3),
        before_sleep=before_sleep_log(logger, logging.DEBUG),
        )
def raise_my_exception():
    global CALL
    CALL += 1
    print(f"CALL = {CALL}, time: {datetime.now()}")
    raise ValueError(f"{CALL} has Value Error")

try:
    raise_my_exception()
except RetryError as e:
    print(f"retry error = {e}")

stats = raise_my_exception.retry.statistics

print(f"Retry 시작 시간: {datetime.now()}")
print(f"시도 횟수: {stats['attempt_number']}")
print(f"함수 첫번째 실행 후 총 경과 시간: {stats['delay_since_first_attempt']}")

""" output
CALL = 1, time: 2024-10-17 15:57:14.372598
DEBUG:__main__:Finished call to '__main__.raise_my_exception' after 0.000(s), this was the 1st time calling it.
DEBUG:__main__:Retrying __main__.raise_my_exception in 0.0 seconds as it raised ValueError: 1 has Value Error.
CALL = 2, time: 2024-10-17 15:57:14.373312
DEBUG:__main__:Finished call to '__main__.raise_my_exception' after 0.001(s), this was the 2nd time calling it.
DEBUG:__main__:Retrying __main__.raise_my_exception in 0.0 seconds as it raised ValueError: 2 has Value Error.
CALL = 3, time: 2024-10-17 15:57:14.373386
DEBUG:__main__:Finished call to '__main__.raise_my_exception' after 0.001(s), this was the 3rd time calling it.
retry error = RetryError[<Future at 0x7fa660774970 state=finished raised ValueError>]
"""

함수가 수행 되고, after_log가 찍히고 before_sleep_log가 찍히는 것을 확인할 수 있습니다.

하지만 이렇게 기본적인 로깅만 구성 하는 것은 추후 로그 분석 할 때 개발자를 피곤하게 만들테니 늘 그러하듯 커스텀 로깅을 구축 하여 가독성 있는 텍스트로 바꾸는 작업이 필요합니다.

import logging
import sys

from tenacity import RetryError, stop_after_attempt, retry, before_sleep_log

logging.basicConfig(
    format="[%(levelname)s|%(filename)s:%(lineno)s] %(asctime)s > %(message)s",
    datefmt='%m/%d/%Y %I:%M:%S',
    level=logging.INFO)

logger = logging.getLogger(__name__)

def my_before_sleep(retry_state):
    if retry_state.attempt_number < 1:
        loglevel = logging.INFO
    else:
        loglevel = logging.WARNING

    logger.log(
        loglevel,
        (f"[{retry_state.fn.__name__}] Retrying ... attempt {retry_state.attempt_number} in {retry_state.next_action.sleep} seconds, result: {retry_state.outcome.exception()}")
    )

CALL = 0

@retry(
        stop=stop_after_attempt(3),
        before_sleep=before_sleep_log(logger, logging.DEBUG),
        )
def raise_my_exception():
    global CALL
    CALL += 1
    print(f"CALL = {CALL}, time: {datetime.now()}")
    raise ValueError(f"{CALL} has Value Error")

try:
    raise_my_exception()
except RetryError as e:
    print(f"retry error = {e}")

stats = raise_my_exception.retry.statistics

print(f"Retry 시작 시간: {datetime.now()}")
print(f"시도 횟수: {stats['attempt_number']}")
print(f"함수 첫번째 실행 후 총 경과 시간: {stats['delay_since_first_attempt']}")

""" output
CALL = 1, time: 2024-10-17 16:41:02.779005
[WARNING|error_list.py:62] 10/17/2024 04:41:02 > [raise_my_exception] Retrying ... attempt 1 in 0.0 seconds, result: 1 has Value Error
CALL = 2, time: 2024-10-17 16:41:02.779206
[WARNING|error_list.py:62] 10/17/2024 04:41:02 > [raise_my_exception] Retrying ... attempt 2 in 0.0 seconds, result: 2 has Value Error
CALL = 3, time: 2024-10-17 16:41:02.779270
retry error = RetryError[<Future at 0x7fb4f8213400 state=finished raised ValueError>]
Retry 시작 시간: 2024-10-17 16:41:02.779304
시도 횟수: 3
함수 첫번째 실행 후 총 경과 시간: 0.000279957996099256
"""

In the real world

실무 프로젝트에서 API 요청이 실패 하는 케이스의 대부분은 재시도 했을 때 정상 응답을 받도록 안내 하는 에러 케이스로 이뤄져 있었습니다.

이로 인해 Retry 정책을 이용 하여 데이터 수집 과정에서 실패 하는 케이스를 다시 시도 하는 운영 방식을 사용 하게 되었습니다.

먼저, API 통신을 유지하던 과정에서 발생한 에러 케이스에 대한 명확한 Raise 처리가 필요합니다. 특정 Exception에 해당 했을 때 Retry 정책이 수행 되어야 하기 때문이죠.

아래 예시를 살펴보면 APIResponseCodeError(APIException)는 API 결과 값에 대한 비정상 상태 코드를 응답 받았을 때 입니다. 그리고, API 결과 값이 정상적으로 변환되지 않았을 때 큰 분류에서 에러를 잡아냅니다.

두 에러에 해당 하는 Exception은 APIException에 해당되기 때문에 이 케이스에 대해서 재수행을 수행 하도록 설정 하면 됩니다.

⚙️ 환경
- Python 3.7.13
- Tenacity 8.2.3

src/example_controller.py

response: Response = requests.post(api_endpoint, headers=header, data=data)
elapsed_time: float = float(f"{(time.perf_counter() - request_time):.3f}")

try:
    api_response: ApiResponse = ApiResponse(response)
    api_request_result_log: DICT = self.service.save_api_request_log(request, api_response, elapsed_time)
    if self._is_api_response_error(api_response):
        raise APIResponseCodeError(
            f"API Response error!, "
            f"response: {api_response.to_dict()}"
        )

except APIResponseDecodeError:
    raise APIException(
        f"API JSON decoded error!, "
        f"\nresponse text: {response.text}")

Retry 모듈을 아래와 같이 구성 하여 필요한 내용을 커스텀 합니다.

src/utils/retries.py

import logging

from tenacity import (
    stop_after_attempt,
    wait_exponential,
    retry_if_exception_type,
    RetryCallState,
    Retrying
)

from shared.exceptions import APIException

TOTAL_RETRIES = 4  # NOTE: 실패 후 총 3번 재실행
WAIT_MULTIPLIER = 2
WAIT_MIN = 2
WAIT_MAX = 10

logger = logging.getLogger(__name__)

def before_callback_retry_messaage(retry_state: RetryCallState):
    """재시도 전 실행되는 콜백 함수"""

    exception = retry_state.outcome.exception()
    logger.warning(
        f"\t * Retry attempt of {retry_state.attempt_number} failed *"
        f"> {retry_state.next_action.sleep} retrying seconds after...\n"
        f"\t * error: {str(exception)} "
    )

def after_callback_finally_message(retry_state: RetryCallState):
    """재시도 후 실행되는 콜백 함수"""

    if retry_state.outcome.failed:
        if retry_state.attempt_number >= TOTAL_RETRIES:
            logger.error(f"\t * Finally failed. error: {str(retry_state.outcome.exception())}")

def api_retry(func):
    max_attempts: int = TOTAL_RETRIES
    wait_min: float = WAIT_MIN
    wait_max: float = WAIT_MAX

    def wrapper(*args, **kwargs):
        retryer = Retrying(
            stop=stop_after_attempt(max_attempts),
            wait=wait_exponential(multiplier=WAIT_MULTIPLIER, min=wait_min, max=wait_max),
            retry=retry_if_exception_type(APIException),
            before_sleep=before_callback_retry_messaage,
            after=after_callback_finally_message,
            reraise=True
        )
        return retryer(func, *args, **kwargs)
    return wrapper

재수행 시 얼마 만큼의 시간을 두고 재수행 할 것인가에 대한 정의가 필요한데, 일정 시간 만큼 고정 시킬 수 있지만 증가 시킬 수도 있습니다. 저는 네트워크 통신을 해야 하기 때문에 일정 시간을 증가 시키는 방법을 이용 하였습니다.

위에서 Multiplier는 초기 값을 설정 한 뒤 multiplier * (2 ** (attempt — 1) 규칙을 따르게 됩니다.

attempt number는 Total 시도 N — 원래 시도 1번 을 통해 재시도가 이뤄집니다. 만약, stop_after_attempt(4) 를 설정 했다면 재시도는 총 4–1 로 3번을 재시도 하게 됩니다.

그럼 최종적으로 해당 데코레이터는 이러한 결과를 나타내게 됩니다.

총 3번을 재수행 하며 수행 될 때 마다 2초 , 4초 , 8초 의 시간 만큼 기다린 뒤 함수를 다시 수행 하게 될 것입니다.

reraise는 본문 내용 위에서 다룬 바 있듯이 RetryException을 기본적으로 잡고 있지만 커스텀 Exception을 사용 하기 위해서는 해당 플래그를 적용 해야 합니다.

최종적으로 아래와 같습니다. 위에서 보았던 컨트롤러 객체에서 사용하는 코드에서 커스텀 데코레이터를 이용 하면 의도한 재수행 정책에 따라 수행 됩니다.

@api_retry
def api_request():
    try:
        api_response: ApiResponse = ApiResponse(response)
        api_request_result_log: DICT = self.service.save_api_request_log(request, api_response, elapsed_time)
        if self._is_api_response_error(api_response):
            raise APIResponseCodeError(
                f"API Response error!, "
                f"response: {api_response.to_dict()}"
            )
    
    except APIResponseDecodeError:
        raise APIException(
            f"API JSON decoded error!, "
            f"\nresponse text: {response.text}")

        ...

이렇게 재수행 정책을 적용한 뒤 실제 API 응답이 실패 -> 성공(Retry) 로 정상적인 응답을 받게 되는 경우가 상당히 증가 하였습니다.

실패 -> 성공 예시

실무 프로젝트 내에는 단순 재시도가 필요 해서 사용 했지만 Airflow처럼 다른 오픈소스를 살펴 보면 대부분 Database Connection에서 많이 사용 되는 것을 확인 할 수 있습니다.

요구사항이 그저 단순 재시도가 필요한 것이라면 Tenacity는 꽤 괜찮은 라이브러리라고 생각 됩니다.

감사합니다.

참고 문서