Skip to content

Commit d92e1b9

Browse files
committed
fix: resolve infinite retry and ffmpeg blocking issues
- Delete SQS message IMMEDIATELY upon receipt (at-most-once delivery) - Redirect ffmpeg stderr to DEVNULL to prevent pipe buffer blocking - Implement application-level retry using send_message with customRetryCount - This prevents SQS visibility timeout issues and zombie processes
1 parent 446a7c9 commit d92e1b9

3 files changed

Lines changed: 42 additions & 11 deletions

File tree

echoshot_ai_server/core/sqs_client.py

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -133,6 +133,20 @@ def __init__(self):
133133
self.queue_url = settings.SQS_QUEUE_URL
134134
self.sqs_client = boto3.client('sqs', region_name=settings.AWS_REGION)
135135

136+
def send_message(self, message_body: Dict[str, Any], delay_seconds: int = 0) -> bool:
137+
"""메시지 전송 (재시도용)"""
138+
try:
139+
self.sqs_client.send_message(
140+
QueueUrl=self.queue_url,
141+
MessageBody=json.dumps(message_body),
142+
DelaySeconds=delay_seconds
143+
)
144+
logger.info(f"Message sent to SQS queue (retry): job_id={message_body.get('job_id') or message_body.get('jobId')}")
145+
return True
146+
except ClientError as e:
147+
logger.error(f"Failed to send message to SQS: {e}")
148+
return False
149+
136150
def receive_messages(self, max_messages: int = 1,
137151
visibility_timeout: int = 300) -> List[Job]:
138152
"""SQS 메시지 수신 및 Job 객체로 변환"""
@@ -176,14 +190,20 @@ def receive_messages(self, max_messages: int = 1,
176190
parsed = parse_sqs_message_body(body)
177191

178192
# Job 객체 생성
193+
# 재시도 횟수: SQS 속성 또는 메시지 본문의 customRetryCount 사용
194+
sqs_retry_count = receive_count - 1 # receive_count는 1부터 시작
195+
custom_retry_count = body.get('customRetryCount', 0)
196+
final_retry_count = max(sqs_retry_count, custom_retry_count)
197+
179198
job = Job(
180199
job_id=parsed['job_id'],
181200
user_id=parsed['user_id'],
182201
task_type=parsed['task_type'],
183202
source_s3_key=parsed['source_s3_key'],
184203
parameters=parsed['parameters'],
185204
receipt_handle=msg['ReceiptHandle'],
186-
metadata=parsed.get('metadata')
205+
metadata=parsed.get('metadata'),
206+
retry_count=final_retry_count
187207
)
188208
jobs.append(job)
189209
logger.debug(f"Successfully converted message to Job: {job.job_id} (task_type={job.task_type}, user_id={job.user_id}, s3_key={parsed['source_s3_key']})")

echoshot_ai_server/services/worker_pool.py

Lines changed: 18 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -89,17 +89,28 @@ def _worker_loop(self, worker_id: int) -> None:
8989

9090
logger.info(f"Worker {worker_id} processing job {job.job_id}")
9191

92-
# Job 처리
92+
# 1. 즉시 삭제 (At-most-once delivery 보장, 중복 실행 방지)
93+
self.sqs_client.delete_message(job.receipt_handle)
94+
logger.info(f"Job {job.job_id} message deleted from SQS immediately to prevent duplicate processing")
95+
96+
# 2. Job 처리
9397
result = self.job_processor.process_job(job)
9498

95-
# SQS 메시지 삭제 (성공/실패 관계없이 항상 삭제)
96-
# 재시도는 SQS의 ApproximateReceiveCount로 관리됨
97-
deleted = self.sqs_client.delete_message(job.receipt_handle)
99+
# 3. 실패 시 재시도 처리 (Application-level retry)
100+
if result.status != JobStatus.COMPLETED:
101+
if job.retry_count < self.job_processor.max_retries:
102+
# 재시도 횟수 증가하여 새 메시지 발행
103+
retry_payload = job.to_dict()
104+
retry_payload['customRetryCount'] = job.retry_count + 1
105+
106+
logger.warning(f"Job {job.job_id} failed. Re-queueing for retry ({retry_payload['customRetryCount']}/{self.job_processor.max_retries})")
107+
self.sqs_client.send_message(retry_payload, delay_seconds=60) # 60초 딜레이
108+
else:
109+
logger.error(f"Job {job.job_id} exceeded max retries ({job.retry_count}). Dropping message.")
98110

111+
# 성공 시에는 이미 삭제했으므로 추가 동작 없음
99112
if result.status == JobStatus.COMPLETED:
100-
logger.info(f"Job {job.job_id} completed successfully, SQS message deleted={deleted}")
101-
else:
102-
logger.warning(f"Job {job.job_id} failed, SQS message deleted={deleted}")
113+
logger.info(f"Job {job.job_id} completed successfully")
103114

104115
except Empty:
105116
continue

echoshot_ai_server/tasks/upscale_task.py

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -251,7 +251,7 @@ def _process_video(
251251
ffmpeg_cmd,
252252
stdin=subprocess.PIPE,
253253
stdout=subprocess.DEVNULL,
254-
stderr=subprocess.PIPE
254+
stderr=subprocess.DEVNULL # 블로킹 방지를 위해 DEVNULL 사용
255255
)
256256

257257
# 프레임 처리
@@ -289,8 +289,8 @@ def _process_video(
289289
ffmpeg_proc.wait()
290290

291291
if ffmpeg_proc.returncode != 0:
292-
stderr = ffmpeg_proc.stderr.read().decode()
293-
raise RuntimeError(f"ffmpeg 오류: {stderr}")
292+
# stderr를 DEVNULL로 보냈으므로 상세 에러 메시지는 확인 불가
293+
raise RuntimeError(f"ffmpeg 오류 발생 (return code: {ffmpeg_proc.returncode})")
294294

295295
def _build_ffmpeg_cmd(
296296
self,

0 commit comments

Comments
 (0)