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은 다음과 같은 중요한 특성을 가지고 있습니다.
- 컨텍스트 관리자(Context Manager)로서의 역할:
async with asyncio.TaskGroup() as tg:구문을 사용하여 생성됩니다. 이async with블록이 종료될 때,TaskGroup은 자신이 생성한 모든 하위 태스크들이 완료될 때까지 기다립니다. - 유한한 작업 집합 관리:
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를 활용하여 안정적이고 효율적인 비동기 애플리케이션을 구축하기 위한 몇 가지 추가 팁입니다.
TaskGroup의 올바른 사용:asyncio.TaskGroup은 주로 명확한 시작과 끝이 있는, 유한한 수의 작업을 병렬로 실행하고 그 결과를 수집하거나 예외를 처리할 때 사용해야 합니다. 예를 들어, 웹 요청을 받아 여러 API 호출을 병렬로 수행하고 결과를 조합하는 경우에 적합합니다. 무한 루프 내에서 지속적으로 태스크를 생성하는 경우에는 적합하지 않습니다.- 태스크 생명 주기 관리:
asyncio.create_task()로 생성된 태스크는 명시적으로await되거나,cancel()호출 후await되어야 합니다. 그렇지 않으면 예상치 못한 시점에 종료되거나, 애플리케이션 종료 시 정리되지 않을 수 있습니다.TaskGroup이나asyncio.gather는 이런 관리를 용이하게 해줍니다. asyncio.Queue의 활용: 지속적으로 들어오는 작업을 처리해야 하는 Producer-Consumer 패턴에서는asyncio.Queue가 강력한 도구입니다. 이를 통해 생산자와 소비자의 속도 차이를 조절하고, 시스템의 안정성을 높일 수 있습니다.- 예외 처리: 비동기 태스크에서 발생하는 예외는 해당 태스크가
await되지 않으면 프로그램 전체를 종료시키지 않을 수 있습니다 (그러나 경고가 발생). 중요한 태스크는 반드시try...except블록으로 감싸 예외를 처리하거나,asyncio.gather(..., return_exceptions=True)와 같이 예외를 반환하도록 처리해야 합니다. - graceful shutdown 구현: 애플리케이션이 종료될 때, 모든 활성 태스크들이 안전하게 정리되고 종료되도록
SIGTERM,SIGINT와 같은 시그널 핸들러를 구현하는 것이 중요합니다. 이 핸들러 내에서 모든 태스크를 취소하고await하는 로직을 포함해야 합니다.
이러한 원칙들을 준수함으로써, Python asyncio 기반의 Redis PUB/SUB 애플리케이션을 더욱 견고하고 효율적으로 운영할 수 있을 것입니다.
![[에러 해결] CUDA 11.1 설치 불가: PyTorch 프로젝트 환경 설정 버전 불일치 원인과 해결 방법 [에러 해결] CUDA 11.1 설치 불가: PyTorch 프로젝트 환경 설정 버전 불일치 원인과 해결 방법](https://dev-error.com/wp-content/plugins/contextual-related-posts/default.png)