JobDri의 비동기 AI 작업을 처리하는 Python 워커 서버입니다.
이 레포는 Spring Boot 메인 서버와 RabbitMQ 사이에서 동작하며, 아래 두 종류의 작업을 비동기로 처리합니다.
- 채용 공고 정리(
JOB_POSTING_INGEST) - 자소서 분석(
ANALYSIS)
작업을 큐에서 구독한 뒤 OpenAI 호출을 수행하고, 결과를 Spring 내부 API로 저장/완료 처리하는 것이 이 레포의 핵심 역할입니다.
메인 백엔드(Spring)는 사용자 요청을 받은 뒤 직접 긴 AI 작업을 수행하지 않습니다. 대신 RabbitMQ에 작업 메시지를 적재하고, 이 워커가 메시지를 소비합니다.
워커는 다음 순서로 동작합니다.
- RabbitMQ 큐를 구독합니다.
- 메시지를 역직렬화해 작업 타입을 판별합니다.
- Spring 내부 API에 작업 상태를
RUNNING으로 반영합니다. - OpenAI를 호출해 채용 공고 추출/분류/생성 또는 자소서 분석을 수행합니다.
- 결과를 Spring 내부 API에 저장합니다.
- 최종 완료 콜백을 전달합니다.
- 실패 시 재시도, DLQ 적재, recovery spool 복구를 수행합니다.
flowchart LR
A["Client"] --> B["Spring Boot API"]
B --> C["RabbitMQ Exchange<br/>jobdri.worker.exchange"]
C --> D["jobdri.job-posting.ingest"]
C --> E["jobdri.analysis.execute"]
D --> F["Python Worker"]
E --> F
F --> G["OpenAI API"]
F --> H["Spring Internal Worker API"]
F --> I["Recovery Spool"]
F --> J["DLQ"]
H --> B
큐 적재는 이 레포가 아니라 메인 백엔드(Spring)에서 수행합니다.
- 채용 공고 작업
APP_WORKER_JOB_POSTING_EXCHANGE+APP_WORKER_JOB_POSTING_ROUTING_KEY - 자소서 분석 작업
APP_WORKER_ANALYSIS_EXCHANGE+APP_WORKER_ANALYSIS_ROUTING_KEY
워커는 적재된 메시지를 소비하는 쪽입니다.
워커 시작 시 app/consumer.py 의 RabbitMqConsumer.start()가 실행되고, 내부 async runtime이 app/async_runtime.py 를 통해 다음 두 큐를 비동기로 구독합니다.
APP_WORKER_JOB_POSTING_QUEUEAPP_WORKER_ANALYSIS_QUEUE
RabbitMQ 연결은 aio-pika 기반이며, 채널 QoS prefetch_count=WORKER_PREFETCH_COUNT 를 유지합니다. 즉, MQ fetch 자체와 내부 HTTP/OpenAI I/O가 같은 event loop에서 겹쳐 처리될 수 있습니다.
또한 task type별 bounded concurrency를 함께 사용합니다.
WORKER_DEFAULT_CONCURRENCY_LIMITWORKER_ANALYSIS_CONCURRENCY_LIMITWORKER_JOB_POSTING_CONCURRENCY_LIMIT
각 메시지는 공통 consumer가 task type별 processor로 dispatch하고, 실제 동시 처리 슬롯은 app/concurrency.py 의 limiter가 제어합니다. 이 구조 덕분에 analysis 적체가 jobposting 전체를 막지 않도록 제한값을 분리해 운영할 수 있습니다.
있습니다. 이 워커는 소비자이면서 일부 상황에서는 다시 RabbitMQ에 메시지를 발행합니다.
- 재시도 시
원본 메시지의
retryCount를 증가시켜 같은 exchange/routing key로 재발행 - 최종 실패 시 DLQ 큐로 발행
즉, 정상 처리 결과는 RabbitMQ로 다시 보내지 않고 Spring 내부 API로 전달하고, 큐 재발행은 재시도/실패 처리에만 사용합니다.
taskType=JOB_POSTING_INGEST
- Spring이 작업 메시지를 큐에 적재합니다.
- 워커가 메시지를 소비합니다.
/api/internal/worker/job-postings/tasks/{taskId}/running으로 상태를RUNNING처리합니다./api/internal/worker/job-postings/ingest/context로 이미지 URL 등 컨텍스트를 조회합니다.- OpenAI로 공고 추출을 수행합니다.
/api/internal/worker/job-postings/classification/candidates로 분류 후보를 조회합니다.- OpenAI로 후보 중 소분류를 선택합니다.
- 신뢰도가 낮으면 저장 없이
/complete로 종료합니다. - 신뢰도가 충분하면 OpenAI로 저장용 공고 내용을 생성합니다.
/result로 중간 결과를 저장합니다./ingest/finalize로 최종 저장/완료를 요청합니다.
taskType=ANALYSIS
- Spring이 분석 작업 메시지를 큐에 적재합니다.
- 워커가 메시지를 소비합니다.
- 큐 대기 시간이
APP_WORKER_ANALYSIS_QUEUE_TIMEOUT_MILLIS를 초과했는지 먼저 검사합니다. /api/internal/worker/analysis/tasks/{taskId}/running으로 상태를RUNNING처리합니다./api/internal/worker/analysis/context로 분석에 필요한 공고/문항/답변 데이터를 조회합니다.- OpenAI로 분석을 수행합니다.
/result로 분석 결과를 저장합니다./complete로 최종 완료 처리합니다.
{
"messageId": "msg-001",
"taskType": "JOB_POSTING_INGEST",
"taskId": "task-job-001",
"userId": 1,
"rawText": "채용 공고 원문",
"imageObjectKey": null,
"retryCount": 0,
"maxRetryCount": 3,
"submittedAt": "2026-07-20T09:00:00Z"
}{
"messageId": "msg-002",
"taskType": "ANALYSIS",
"taskId": "task-analysis-001",
"userId": 1,
"mockApplyId": 42,
"retryCount": 0,
"maxRetryCount": 3,
"submittedAt": "2026-07-20T09:10:00Z"
}스키마 원본은 app/schemas.py 에 정의되어 있습니다.
이 워커는 단순 consume 후 종료하지 않고, 실패 내성을 갖춘 구조로 작성되어 있습니다.
RetryableWorkerError발생 시 재시도 대상으로 간주합니다.retryCount + 1값을 담아 같은 큐 토폴로지로 재발행합니다.- 최대 재시도 횟수는 작업별로 다릅니다.
- 채용 공고:
WORKER_MAX_RETRY_COUNT - 자소서 분석:
APP_WORKER_ANALYSIS_MAX_RETRY_COUNT
- 채용 공고:
아래 경우에는 DLQ로 보냅니다.
- 비재시도 오류(
NonRetryableWorkerError) - 재시도 횟수 초과
설정 키는 다음과 같습니다.
APP_WORKER_JOB_POSTING_DLQAPP_WORKER_ANALYSIS_DLQ
OpenAI 호출은 성공했지만 Spring 내부 API로 최종 완료 콜백을 보내는 단계에서 실패할 수 있습니다. 이때 결과를 잃지 않도록 파일 기반 spool에 pending delivery를 남깁니다.
- 기본 경로
APP_WORKER_RECOVERY_SPOOL_DIR=.worker-spool - terminal message ledger 경로
APP_WORKER_TERMINAL_MESSAGE_DIR=.worker-spool/terminal-messages
워커는 시작 시점과 주기적 백그라운드 루프에서 spool 파일을 다시 읽어 미전달 완료 요청을 재전송합니다.
관련 구현은 다음 파일에 있습니다.
이 워커는 Spring 공개 API가 아니라 내부 worker API를 호출합니다. 공통 헤더로 X-Internal-Api-Key를 사용합니다.
주요 연동 엔드포인트는 다음과 같습니다.
POST /api/internal/worker/job-postings/tasks/{taskId}/runningPOST /api/internal/worker/job-postings/ingest/contextPOST /api/internal/worker/job-postings/classification/candidatesPOST /api/internal/worker/job-postings/tasks/{taskId}/resultPOST /api/internal/worker/job-postings/ingest/finalizePOST /api/internal/worker/job-postings/tasks/{taskId}/retryPOST /api/internal/worker/job-postings/tasks/{taskId}/failedGET /api/internal/worker/job-postings/tasks/{taskId}
POST /api/internal/worker/analysis/tasks/{taskId}/runningPOST /api/internal/worker/analysis/contextPOST /api/internal/worker/analysis/tasks/{taskId}/resultPOST /api/internal/worker/analysis/tasks/{taskId}/completePOST /api/internal/worker/analysis/tasks/{taskId}/retryPOST /api/internal/worker/analysis/tasks/{taskId}/failedGET /api/internal/worker/analysis/tasks/{taskId}
구현은 app/api_client.py 에 있으며, 현재는 sync 메서드와 httpx.AsyncClient 기반 async 메서드를 함께 제공합니다.
app/openai_client.py 의 JobPostingOpenAiWorker 가 아래 3단계를 담당합니다.
extract공고 텍스트 또는 이미지에서 구조화 정보 추출classifySpring이 준 분류 후보 중 가장 적합한 소분류 선택generate저장 가능한 형태의 정제된 공고 내용 생성
AnalysisOpenAiWorker 가 문항별 답변과 공고 맥락을 바탕으로 종합 점수와 피드백을 생성합니다.
현재 OpenAI 호출도 sync/async wrapper를 모두 가지고 있으며, 실제 consumer 경로는 async client를 사용합니다.
기본 모델은 아래 환경변수로 제어됩니다.
OPENAI_JOB_POSTING_MODEL=gpt-4o-miniOPENAI_ANALYSIS_MODEL=gpt-4.1-mini
app/
main.py # FastAPI 앱과 health endpoint
worker.py # CLI worker 진입점
consumer.py # worker 조립과 공통 consume helper
async_runtime.py # aio-pika 기반 async consumer runtime
concurrency.py # task type별 bounded concurrency limiter
processors.py # task type별 processor 분리
delivery.py # result store / finalize / complete 전달 계층
api_client.py # Spring 내부 API 클라이언트 (sync + async)
openai_client.py # OpenAI 호출 로직 (sync + async)
async_utils.py # sync/async bridge helper
recovery.py # recovery spool 저장/복구
schemas.py # 메시지/응답 스키마
logging_utils.py # 구조화 로그 필터
tests/
test_prometheus_metrics.py
test_recovery_flow.py
deploy/
docker-compose.worker.prod.yml
docs/
BACKEND_SERVER_DEPLOY.md
RENDER_DEPLOY_CHECKLIST.md
- Python 3.12 이상 권장
- RabbitMQ 접근 가능
- Spring 내부 worker API 접근 가능
- OpenAI API 키
이 레포는 Dockerfile 기준으로 Python 3.12 환경에서 실행됩니다.
python3 -m venv .venv
source .venv/bin/activate
pip install -r requirements.txt최소한 아래 값은 필요합니다.
APP_WORKER_INTERNAL_API_KEY=change-me
SPRING_API_BASE_URL=http://localhost:8080
OPENAI_API_KEY=sk-...
RABBITMQ_HOST=localhost
RABBITMQ_PORT=5672
RABBITMQ_USERNAME=guest
RABBITMQ_PASSWORD=guest
RABBITMQ_VHOST=/
APP_WORKER_JOB_POSTING_EXCHANGE=jobdri.worker.exchange
APP_WORKER_JOB_POSTING_QUEUE=jobdri.job-posting.ingest
APP_WORKER_JOB_POSTING_ROUTING_KEY=job-posting.ingest
APP_WORKER_JOB_POSTING_DLQ=jobdri.job-posting.ingest.dlq
APP_WORKER_ANALYSIS_EXCHANGE=jobdri.worker.exchange
APP_WORKER_ANALYSIS_QUEUE=jobdri.analysis.execute
APP_WORKER_ANALYSIS_ROUTING_KEY=analysis.execute
APP_WORKER_ANALYSIS_DLQ=jobdri.analysis.execute.dlq
WORKER_DEFAULT_CONCURRENCY_LIMIT=1
WORKER_ANALYSIS_CONCURRENCY_LIMIT=1
WORKER_JOB_POSTING_CONCURRENCY_LIMIT=1전체 목록은 app/config.py 와 deploy/docker-compose.worker.prod.yml 를 참고하면 됩니다.
uvicorn app.main:app --host 0.0.0.0 --port 8000- health endpoint
GET /health - metrics endpoint
GET /metrics
FastAPI 앱 startup에서 RabbitMQ consumer도 함께 시작되므로, 별도 백그라운드 워커 프로세스를 추가로 띄우지 않습니다.
이미지는 Dockerfile 기준으로 빌드되며 기본 실행 명령은 아래와 같습니다.
CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"]운영 배포 예시는 deploy/docker-compose.worker.prod.yml 에 있습니다.
상세 운영 문서는 아래 파일을 참고하세요.
로그에는 기본적으로 아래 컨텍스트가 함께 남습니다.
taskIdmessageIdworkerIdretryCount
정상 기동 시 확인할 대표 로그:
Worker process started.RabbitMQ consumer started.
기존 Prometheus 메트릭은 유지하면서, 2차 구현에서는 async 경로와 운영 병목을 더 직접 해석할 수 있도록 아래를 정리했습니다.
worker_task_queue_wait_duration_seconds{task_type}큐 적체 시간을 봅니다. 현재 병목 1순위 확인용입니다.worker_task_processing_duration_seconds{task_type,outcome}worker 내부 end-to-end 처리 시간을 봅니다.llm_request_duration_seconds{task_type,operation,outcome}OpenAI 호출 latency를 봅니다.worker_internal_api_duration_seconds{task_type,endpoint,method,outcome}Spring internal API 전체 호출 시간을 봅니다.worker_context_fetch_duration_seconds{task_type,endpoint,outcome}context fetch 전용 뷰입니다.worker_callback_duration_seconds{task_type,endpoint,outcome}complete/finalize/retry/failed callback 전용 뷰입니다.worker_task_inflight{task_type}현재 동시에 처리 중인 task 수입니다.worker_task_concurrency_limit{task_type}task type별 설정 상한입니다.inflight와 같이 보면 슬롯 포화 여부를 해석하기 쉽습니다.worker_task_retry_count_total{task_type,reason}retry 전환량입니다.
- analysis 응답 schema 검증 실패는
VALIDATION_ERROR+ non-retryable로 분류합니다. - job posting
classify/generate의 schema 검증 실패는 fallback으로 처리하고,llm_request_duration_seconds{outcome="fallback"}로 기록합니다. - fallback도 성공으로 섞지 않고 별도 outcome으로 유지해, success latency와 fallback latency를 분리해서 볼 수 있게 했습니다.
- retryable / non-retryable 의미는 기존 비즈니스 계약을 유지합니다.
RATE_LIMIT,OPENAI_TIMEOUT, 일반INTERNAL_ERROR, Spring API transport/server error는 재시도 가능- validation/parsing failure와 4xx 성격 오류는 비재시도
큐 대기 p95:
histogram_quantile(
0.95,
sum by (le, task_type) (
rate(worker_task_queue_wait_duration_seconds_bucket[5m])
)
)
큐 대기 p99:
histogram_quantile(
0.99,
sum by (le, task_type) (
rate(worker_task_queue_wait_duration_seconds_bucket[5m])
)
)
처리 시간 p95:
histogram_quantile(
0.95,
sum by (le, task_type, outcome) (
rate(worker_task_processing_duration_seconds_bucket[5m])
)
)
LLM 요청 p95:
histogram_quantile(
0.95,
sum by (le, task_type, operation, outcome) (
rate(llm_request_duration_seconds_bucket[5m])
)
)
inflight 평균:
avg_over_time(worker_task_inflight[5m])
inflight 최대:
max_over_time(worker_task_inflight[5m])
설정 상한 대비 사용량:
worker_task_inflight / clamp_min(worker_task_concurrency_limit, 1)
retry rate:
sum by (task_type, reason) (
rate(worker_task_retry_count_total[5m])
)
context fetch p95:
histogram_quantile(
0.95,
sum by (le, task_type, endpoint, outcome) (
rate(worker_context_fetch_duration_seconds_bucket[5m])
)
)
callback p95:
histogram_quantile(
0.95,
sum by (le, task_type, endpoint, outcome) (
rate(worker_callback_duration_seconds_bucket[5m])
)
)
python3 -m unittest \
tests/test_prometheus_metrics.py \
tests/test_worker_concurrency.py \
tests/test_recovery_flow.py \
tests/test_observability_logging.py현재 테스트는 아래를 검증합니다.
- recovery spool 기반 재전송과 완료 콜백 복구
- validation error / fallback / retry 메트릭 분류
- task type별 inflight / concurrency metric
- observability 로그 필드 정합성