1. 워커 서비스의 신뢰성 문제 분석
기존 아키텍처에서 워커(Worker)가 대규모 언어 모델(LLM) API를 호출할 때 블로킹(Blocking) 대기 상태가 발생합니다. 이 과정에서 워커 컨테이너가 예기치 않게 종료되거나 중단될 경우, 해당 작업은 데이터베이스상에서 계속 'processing' 상태로 남게 됩니다. Redis의 작업 큐에서 이미 메시지가 소비되었기 때문에, 시스템은 해당 작업을 다시 식별하거나 재시도할 수 없는 결함이 발생합니다.
이러한 문제를 해결하기 위해 '순찰(Inspection) 서비스'를 도입하여 타임아웃이 발생한 'processing' 상태의 작업을 정기적으로 복구해야 합니다. 성능 최적화를 위해 MySQL 전체를 조회하는 대신, Redis에 별도의 processing 큐와 processing_ts(시작 시간 기록용) 해시 구조를 구축합니다. 순찰 서비스는 이 Redis 데이터를 기반으로 비정상 작업을 감지하고 복구합니다.
워커 서비스는 작업을 가져오는 동시에 '처리 중' 상태로 옮기는 과정이 중단되지 않도록 원자적 연산을 수행해야 합니다.
2. 로직 구현 및 최적화
2.1 워커의 원자적 작업 이동 및 시간 기록
워커가 작업을 가져올 때 BRPOPLPUSH 명령을 사용하여 작업을 'ready' 큐에서 'processing' 큐로 원자적으로 이동시킵니다. 이후 해시 구조에 시작 시간을 기록합니다.
import json
import time
import logging
# Redis 키 설정
QUEUE_READY = "llm_test:task:ready"
QUEUE_ACTIVE = "llm_test:task:active"
HASH_START_TIME = "llm_test:task:start_at"
def run_worker_node():
"""문서 테스트 워커 메인 루프"""
logging.info("Worker node started. Monitoring task queue...")
while True:
try:
# 원자적으로 작업을 active 큐로 이동
task_raw = redis_conn.brpoplpush(QUEUE_READY, QUEUE_ACTIVE, timeout=15)
if not task_raw:
continue
try:
task_data = json.loads(task_raw.decode("utf-8"))
tid = task_data["task_id"]
except (json.JSONDecodeError, KeyError) as e:
logging.error(f"Payload error: {e}")
redis_conn.lrem(QUEUE_ACTIVE, 1, task_raw)
continue
# 시작 시간 기록
redis_conn.hset(HASH_START_TIME, tid, int(time.time()))
try:
execute_llm_task(tid)
finally:
# 작업 완료 후 큐 및 시간 정보 삭제
redis_conn.lrem(QUEUE_ACTIVE, 1, task_raw)
redis_conn.hdel(HASH_START_TIME, tid)
except Exception as ex:
logging.error(f"Worker loop error: {ex}")
time.sleep(5)
2.2 순찰 서비스(Reaper) 및 데이터베이스 상태 복구
설정된 임계치(예: 600초)를 초과하여 처리 중인 작업은 워커 장애로 간주합니다. 해당 작업은 데이터베이스에서 다시 'pending' 상태로 되돌리고 'ready' 큐에 재진입시켜 다른 워커가 처리할 수 있도록 합니다.
데이터베이스의 원자적 상태 변경을 위해 SQLAlchemy를 활용한 복구 로직을 구현합니다.
def restore_stale_task(task_id: int, limit_time) -> bool:
"""타임아웃된 작업을 다시 대기 상태로 복구"""
with session_scope() as session:
query = (
update(TaskModel).where(
TaskModel.id == task_id,
TaskModel.status == "processing",
TaskModel.started_at < limit_time
).values(
status="pending",
retry_count=TaskModel.retry_count + 1,
started_at=None
)
)
res = session.execute(query)
return res.rowcount == 1
다음은 주기적으로 Redis를 스캔하여 정체된 작업을 식별하고 복구하는 Reaper 루프입니다. 여러 Reaper가 동시에 작동하더라도 데이터베이스의 원자적 업데이트(UPDATE) 조건문 덕분에 중복 복구가 방지됩니다.
def reaper_process():
"""지속적으로 active 큐를 감시하고 고립된 작업 복구"""
timeout_threshold = 600
while True:
try:
now = int(time.time())
active_tasks = redis_conn.lrange(QUEUE_ACTIVE, 0, -1)
for task_raw in active_tasks:
data = json.loads(task_raw.decode("utf-8"))
tid = data.get("task_id")
# 시작 시간 확인
start_ts = redis_conn.hget(HASH_START_TIME, tid)
if not start_ts or (now - int(start_ts)) < timeout_threshold:
continue
logging.warning(f"Task {tid} detected as stale. Reclaiming...")
# DB 상태 업데이트 시도
limit_dt = datetime.utcnow() - timedelta(seconds=timeout_threshold)
if restore_stale_task(tid, limit_dt):
# Redis 정보 정리 및 재입고
redis_conn.lrem(QUEUE_ACTIVE, 1, task_raw)
redis_conn.hdel(HASH_START_TIME, tid)
redis_conn.lpush(QUEUE_READY, task_raw)
logging.info(f"Task {tid} successfully requeued.")
except Exception as e:
logging.error(f"Reaper error: {e}")
time.sleep(30)
3. 검증 테스트
워커가 작업을 처리하는 도중 강제로 프로세스를 종료하는 시나리오를 통해 기능을 검증합니다. 이때 MySQL의 상태는 'processing'이며 Redis의 'ready' 큐에는 해당 태스크가 존재하지 않고 'active' 큐와 'start_at' 해시에만 데이터가 남게 됩니다.
순찰 서비스(Reaper)를 가동하면 로그를 통해 타임아웃이 발생한 특정 작업을 감지하고, 해당 작업을 다시 'ready' 큐로 되돌리는 과정을 확인할 수 있습니다. 결과적으로 재기동된 워커는 복구된 작업을 다시 가져와 처리를 완료하며, 데이터베이스에는 재시도 횟수와 갱신된 시간이 정확히 기록됩니다. 이를 통해 분산 환경에서의 작업 유실 문제를 효과적으로 해결할 수 있음을 확인하였습니다.