[에러 해결] [Python asyncio] Redis PUB/SUB 장기 실행 TaskGroup 사용의 함정과 해결책

Python의 비동기 프로그래밍은 I/O 바운드 작업에서 탁월한 성능을 발휘합니다. 특히 Redis의 PUB/SUB 기능을 활용하여 메시지를 실시간으로 처리할 때, asyncio는 매우 유용한 도구입니다. 하지만 asyncio.TaskGroup을 잘못 사용하면 예상치 못한 문제에 직면할 수 있습니다. 이번 기술 블로그에서는 Redis PUB/SUB 환경에서 asyncio.TaskGroup을 장기 실행 태스크로 활용하려 할 때 발생할 수 있는 문제점과, 이를 효율적이고 안정적으로 해결하는 방법을 심층적으로 분석합니다.

1. 에러 발생 상황

개발자는 Redis PUB/SUB 채널에서 들어오는 메시지를 효율적으로 처리하기 위해 Python의 asyncio를 사용하고 있었습니다. 특히, 여러 메시지를 동시에 비동기적으로 처리하기 위해 asyncio.TaskGroup을 활용하려는 시도를 했습니다.

코드의 핵심 아이디어는 다음과 같습니다:

  • 단일 리스너 태스크(_listen_channels)가 모든 Redis 채널의 메시지를 수신합니다.
  • 이 리스너 태스크 내부에서 asyncio.TaskGroup을 사용하여 수신된 각 메시지를 처리하는 작업을 생성합니다.

다음은 문제가 발생한 코드 예시입니다.

def start_listen_channels(self):
    # create listener task
    if self._listen_channels_task is None:
        self._listen_channels_task = asyncio.create_task(
            self._listen_channels()
        )

async def _listen_channels(self):
    logger.info("start listenning channels")
    # tasks group
    async with asyncio.TaskGroup() as tg:
        # loop messages in pubsub
        async for message in self.pubsub.listen():
            message_type = message.get("type")
            message_data = message.get("data", {})
            # get message channel
            channel: str = message.get("channel")
            channel_id = channel.removeprefix(self.CHANNELS_PREFIX)
            # handle message for channel
            tg.create_task(self.on_channel_message(channel_id, message_data))

이 코드를 구현한 후, 개발자는 다음과 같은 근본적인 질문을 던졌습니다:

  • “이 방식이 과연 효율적일까?”
  • asyncio.TaskGroup이 장기 실행(long-living) 태스크 그룹으로 사용될 수 있을까?”

2. 명확한 발생 원인

이 질문에 대한 핵심은 asyncio.TaskGroup의 동작 방식에 있습니다. asyncio.TaskGroup은 다음과 같은 중요한 특성을 가지고 있습니다.

  1. 컨텍스트 관리자(Context Manager)로서의 역할: async with asyncio.TaskGroup() as tg: 구문을 사용하여 생성됩니다. 이 async with 블록이 종료될 때, TaskGroup은 자신이 생성한 모든 하위 태스크들이 완료될 때까지 기다립니다.
  2. 유한한 작업 집합 관리: TaskGroup은 일반적으로 특정한 일련의 병렬 작업을 시작하고, 그 작업들이 모두 완료될 때까지 기다리거나 예외를 처리하는 데 적합하게 설계되었습니다. 즉, 명확한 시작과 끝이 있는 작업 그룹에 적합합니다.

제시된 코드에서는 async for message in self.pubsub.listen(): 루프가 무기한으로 실행됩니다. Redis PUB/SUB 리스너는 메시지가 계속 들어오는 한 종료되지 않기 때문입니다. 이 무한 루프는 async with asyncio.TaskGroup() as tg: 블록 내부에 위치하므로, TaskGroup 컨텍스트는 절대 종료되지 않습니다.

이것이 문제의 핵심입니다. TaskGroup 컨텍스트가 종료되지 않으므로, tg.create_task(...)로 생성된 메시지 처리 태스크들은 TaskGroup에 의해 정식으로 “완료 대기(await)” 되거나 관리되지 않습니다. 각 on_channel_message 태스크는 독립적으로 실행되지만, TaskGroup의 생명주기와는 무관하게 떠다니게 됩니다. 이는 다음과 같은 잠재적인 문제를 야기할 수 있습니다.

  • 자원 누수 가능성: 태스크가 비정상적으로 종료되거나, 처리 중 예상치 못한 문제가 발생했을 때, TaskGroup이 이들을 모니터링하거나 정리할 기회가 없습니다.
  • 예외 처리 복잡성: 하위 태스크에서 발생하는 예외가 TaskGroup 컨텍스트로 전파되지 않으므로, 예외를 적절히 처리하기 어려워집니다.
  • 종료 로직의 부재: 애플리케이션 종료 시, TaskGroup이 활성 태스크들을 종료하거나 정리하는 데 도움을 주지 못합니다.

결론적으로, asyncio.TaskGroup은 “장기 실행되며 무기한으로 새로운 태스크를 생성하는” 패턴에는 적합하지 않습니다. TaskGroup은 유한한 작업 집합을 효율적으로 병렬 처리하고 결과를 모을 때 가장 빛을 발합니다.

3. 해결 방법 및 코드 예시

Redis PUB/SUB과 같이 지속적으로 메시지를 수신하고 처리해야 하는 환경에서는 asyncio.TaskGroup의 본래 목적에 맞는 다른 패턴을 사용하는 것이 좋습니다. 여기 두 가지 주요 해결 방법을 제안합니다.

3.1. 각 메시지 처리 태스크를 독립적으로 관리 (간단한 경우)

가장 간단한 방법은 각 메시지 처리 작업을 asyncio.create_task()를 사용하여 독립적인 태스크로 생성하는 것입니다. 이 경우 TaskGroup의 오버헤드를 피하고, 태스크들이 서로에게 영향을 미치지 않도록 할 수 있습니다. 하지만 태스크의 생명 주기 관리(취소, 예외 처리 등)는 개발자의 몫이 됩니다.

import asyncio
import logging

logger = logging.getLogger(__name__)

class RedisListener:
    def __init__(self, pubsub_client, channels_prefix="channel:"):
        self.pubsub = pubsub_client  # Redis pubsub 클라이언트 인스턴스
        self.CHANNELS_PREFIX = channels_prefix
        self._listen_channels_task = None
        self._active_tasks = set() # 활성 태스크를 추적하기 위한 집합 (선택 사항)

    async def on_channel_message(self, channel_id, message_data):
        """메시지를 실제 처리하는 비동기 함수"""
        logger.info(f"[{channel_id}] 메시지 처리 시작: {message_data}")
        await asyncio.sleep(1) # 실제 작업 시뮬레이션
        logger.info(f"[{channel_id}] 메시지 처리 완료: {message_data}")
        # 태스크 완료 후 집합에서 제거
        current_task = asyncio.current_task()
        if current_task in self._active_tasks:
            self._active_tasks.remove(current_task)

    async def _listen_channels(self):
        logger.info("채널 리스닝 시작")
        try:
            async for message in self.pubsub.listen():
                if message.get("type") == "message":
                    message_data = message.get("data", {})
                    channel: str = message.get("channel")
                    channel_id = channel.removeprefix(self.CHANNELS_PREFIX)

                    # 각 메시지 처리를 독립적인 태스크로 생성
                    task = asyncio.create_task(self.on_channel_message(channel_id, message_data))
                    # 선택적으로 태스크를 추적하여 종료 시 관리
                    self._active_tasks.add(task)
                    task.add_done_callback(self._active_tasks.discard) # 태스크 완료 시 집합에서 자동 제거
                else:
                    logger.debug(f"수신된 메시지 유형: {message.get('type')}")
        except asyncio.CancelledError:
            logger.info("채널 리스너 태스크 취소됨.")
        except Exception as e:
            logger.error(f"채널 리스너에서 예외 발생: {e}", exc_info=True)
        finally:
            logger.info("채널 리스너 종료 중.")

    def start_listen_channels(self):
        if self._listen_channels_task is None:
            self._listen_channels_task = asyncio.create_task(
                self._listen_channels()
            )
            logger.info("리스너 태스크 생성 완료.")

    async def stop_listen_channels(self):
        if self._listen_channels_task:
            logger.info("리스너 태스크 취소 요청...")
            self._listen_channels_task.cancel()
            await self._listen_channels_task
            logger.info("리스너 태스크 정상 종료.")
        
        # 남아있는 활성 태스크들을 기다리거나 취소
        if self._active_tasks:
            logger.info(f"남아있는 {len(self._active_tasks)}개의 메시지 처리 태스크를 기다립니다...")
            await asyncio.gather(*self._active_tasks, return_exceptions=True)
            logger.info("모든 메시지 처리 태스크 완료.")

# 사용 예시 (실제 Redis 클라이언트가 필요합니다)
# async def main():
#     # 가상의 Redis PubSub 클라이언트 (실제 클라이언트로 교체 필요)
#     class MockPubSub:
#         async def listen(self):
#             for i in range(1, 10):
#                 await asyncio.sleep(0.5)
#                 yield {"type": "message", "channel": "channel:my_channel", "data": {"id": i, "content": f"data {i}"}}
#             # 무한 루프를 시뮬레이션하려면 이 부분을 주석 처리
#             # while True:
#             #     await asyncio.sleep(0.5)
#             #     yield {"type": "message", "channel": "channel:my_channel", "data": {"id": i, "content": f"data {i}"}}
#             #     i += 1

#     listener = RedisListener(MockPubSub())
#     listener.start_listen_channels()
#     
#     await asyncio.sleep(5) # 리스너가 잠시 실행되도록 기다림
#     await listener.stop_listen_channels()
#     logger.info("애플리케이션 종료.")

# if __name__ == "__main__":
#     logging.basicConfig(level=logging.INFO)
#     asyncio.run(main())

3.2. asyncio.Queue와 워커 태스크 패턴 활용 (권장)

더 견고하고 확장 가능한 접근 방식은 `asyncio.Queue`를 사용하는 것입니다. 리스너 태스크는 메시지를 받아 큐에 넣고, 별도의 워커(worker) 태스크들이 큐에서 메시지를 꺼내 처리하는 방식입니다. 이 패턴은 동시성 제어, 백프레셔(backpressure) 관리, 그리고 안정적인 종료 처리에 매우 효과적입니다.

import asyncio
import logging

logger = logging.getLogger(__name__)

class RedisQueueListener:
    def __init__(self, pubsub_client, channels_prefix="channel:", num_workers=5):
        self.pubsub = pubsub_client
        self.CHANNELS_PREFIX = channels_prefix
        self.message_queue = asyncio.Queue()
        self.num_workers = num_workers
        self._listener_task = None
        self._worker_tasks = []

    async def _process_message(self, channel_id, message_data):
        """실제 메시지 처리 로직"""
        logger.info(f"[Worker] [{channel_id}] 메시지 처리 시작: {message_data}")
        await asyncio.sleep(1) # 실제 작업 시뮬레이션
        logger.info(f"[Worker] [{channel_id}] 메시지 처리 완료: {message_data}")

    async def _worker(self, worker_id):
        """큐에서 메시지를 꺼내 처리하는 워커 태스크"""
        logger.info(f"워커 {worker_id} 시작.")
        try:
            while True:
                channel_id, message_data = await self.message_queue.get()
                try:
                    await self._process_message(channel_id, message_data)
                finally:
                    self.message_queue.task_done() # 작업 완료 알림
        except asyncio.CancelledError:
            logger.info(f"워커 {worker_id} 취소됨.")
        except Exception as e:
            logger.error(f"워커 {worker_id} 에서 예외 발생: {e}", exc_info=True)
        finally:
            logger.info(f"워커 {worker_id} 종료.")

    async def _listener(self):
        """Redis PUB/SUB 메시지를 수신하여 큐에 넣는 리스너 태스크"""
        logger.info("Redis 채널 리스닝 시작.")
        try:
            async for message in self.pubsub.listen():
                if message.get("type") == "message":
                    message_data = message.get("data", {})
                    channel: str = message.get("channel")
                    channel_id = channel.removeprefix(self.CHANNELS_PREFIX)
                    await self.message_queue.put((channel_id, message_data)) # 큐에 메시지 추가
                else:
                    logger.debug(f"수신된 메시지 유형: {message.get('type')}")
        except asyncio.CancelledError:
            logger.info("Redis 리스너 태스크 취소됨.")
        except Exception as e:
            logger.error(f"Redis 리스너에서 예외 발생: {e}", exc_info=True)
        finally:
            logger.info("Redis 리스너 종료 중.")

    async def start(self):
        """리스너와 워커 태스크를 시작"""
        logger.info("리스너 및 워커 태스크 시작 준비.")
        self._listener_task = asyncio.create_task(self._listener())
        
        # TaskGroup을 사용하여 워커 태스크들을 관리
        # 워커 태스크들은 무한 루프를 돌지만, 종료 시점에 TaskGroup이 이들을 취소하고 정리할 수 있음
        self._worker_tasks = []
        for i in range(self.num_workers):
            self._worker_tasks.append(asyncio.create_task(self._worker(i+1)))
        
        logger.info(f"리스너 태스크 시작 및 {self.num_workers}개의 워커 태스크 생성 완료.")

    async def stop(self):
        """모든 태스크를 안전하게 종료"""
        logger.info("애플리케이션 종료 요청 받음. 모든 태스크 종료 시작.")
        
        # 1. 리스너 태스크 취소 (더 이상 메시지를 큐에 넣지 않음)
        if self._listener_task:
            self._listener_task.cancel()
            await self._listener_task
            logger.info("리스너 태스크 종료 완료.")

        # 2. 큐에 남아있는 모든 작업이 처리될 때까지 대기
        logger.info(f"큐에 남아있는 작업 ({self.message_queue.qsize()}개) 처리 대기 중...")
        await self.message_queue.join()
        logger.info("큐의 모든 작업 처리 완료.")

        # 3. 워커 태스크들 취소
        for task in self._worker_tasks:
            task.cancel()
        await asyncio.gather(*self._worker_tasks, return_exceptions=True)
        logger.info("모든 워커 태스크 종료 완료.")
        
        logger.info("모든 태스크가 안전하게 종료되었습니다.")

# 사용 예시 (실제 Redis 클라이언트가 필요합니다)
# async def main():
#     # 가상의 Redis PubSub 클라이언트 (실제 클라이언트로 교체 필요)
#     class MockPubSub:
#         async def listen(self):
#             i = 0
#             while True:
#                 await asyncio.sleep(0.1) # 빠르게 메시지 생성
#                 yield {"type": "message", "channel": "channel:my_channel", "data": {"id": i, "content": f"data {i}"}}
#                 i += 1

#     listener_service = RedisQueueListener(MockPubSub(), num_workers=3)
#     await listener_service.start()
#     
#     # Ctrl+C 등으로 종료하기 전까지 실행
#     try:
#         await asyncio.sleep(10) # 서비스 실행 시간
#     except asyncio.CancelledError:
#         pass
#     finally:
#         await listener_service.stop()
#     logger.info("애플리케이션 종료.")

# if __name__ == "__main__":
#     logging.basicConfig(level=logging.INFO)
#     asyncio.run(main())

이 패턴의 장점은 다음과 같습니다:

  • 동시성 제어: num_workers를 통해 동시에 처리될 수 있는 메시지 수를 제한하여 시스템 과부하를 방지할 수 있습니다.
  • 백프레셔: 큐가 가득 차면 put() 호출이 대기하여, 메시지 생산 속도가 처리 속도보다 빠를 때 시스템이 안정적으로 유지됩니다.
  • 안정적인 종료: queue.join()을 사용하여 큐에 있는 모든 메시지가 처리될 때까지 기다린 후 워커 태스크를 종료하여 데이터 손실을 방지할 수 있습니다.

4. 향후 예방을 위한 팁

asyncio를 활용하여 안정적이고 효율적인 비동기 애플리케이션을 구축하기 위한 몇 가지 추가 팁입니다.

  1. TaskGroup의 올바른 사용: asyncio.TaskGroup은 주로 명확한 시작과 끝이 있는, 유한한 수의 작업을 병렬로 실행하고 그 결과를 수집하거나 예외를 처리할 때 사용해야 합니다. 예를 들어, 웹 요청을 받아 여러 API 호출을 병렬로 수행하고 결과를 조합하는 경우에 적합합니다. 무한 루프 내에서 지속적으로 태스크를 생성하는 경우에는 적합하지 않습니다.
  2. 태스크 생명 주기 관리: asyncio.create_task()로 생성된 태스크는 명시적으로 await되거나, cancel() 호출 후 await되어야 합니다. 그렇지 않으면 예상치 못한 시점에 종료되거나, 애플리케이션 종료 시 정리되지 않을 수 있습니다. TaskGroup이나 asyncio.gather는 이런 관리를 용이하게 해줍니다.
  3. asyncio.Queue의 활용: 지속적으로 들어오는 작업을 처리해야 하는 Producer-Consumer 패턴에서는 asyncio.Queue가 강력한 도구입니다. 이를 통해 생산자와 소비자의 속도 차이를 조절하고, 시스템의 안정성을 높일 수 있습니다.
  4. 예외 처리: 비동기 태스크에서 발생하는 예외는 해당 태스크가 await되지 않으면 프로그램 전체를 종료시키지 않을 수 있습니다 (그러나 경고가 발생). 중요한 태스크는 반드시 try...except 블록으로 감싸 예외를 처리하거나, asyncio.gather(..., return_exceptions=True)와 같이 예외를 반환하도록 처리해야 합니다.
  5. graceful shutdown 구현: 애플리케이션이 종료될 때, 모든 활성 태스크들이 안전하게 정리되고 종료되도록 SIGTERM, SIGINT와 같은 시그널 핸들러를 구현하는 것이 중요합니다. 이 핸들러 내에서 모든 태스크를 취소하고 await하는 로직을 포함해야 합니다.

이러한 원칙들을 준수함으로써, Python asyncio 기반의 Redis PUB/SUB 애플리케이션을 더욱 견고하고 효율적으로 운영할 수 있을 것입니다.

댓글 남기기