Skip to content

Latest commit

ย 

History

History
504 lines (372 loc) ยท 16.3 KB

File metadata and controls

504 lines (372 loc) ยท 16.3 KB

JobDri Analysis Server

JobDri์˜ ๋น„๋™๊ธฐ AI ์ž‘์—…์„ ์ฒ˜๋ฆฌํ•˜๋Š” Python ์›Œ์ปค ์„œ๋ฒ„์ž…๋‹ˆ๋‹ค.

์ด ๋ ˆํฌ๋Š” Spring Boot ๋ฉ”์ธ ์„œ๋ฒ„์™€ RabbitMQ ์‚ฌ์ด์—์„œ ๋™์ž‘ํ•˜๋ฉฐ, ์•„๋ž˜ ๋‘ ์ข…๋ฅ˜์˜ ์ž‘์—…์„ ๋น„๋™๊ธฐ๋กœ ์ฒ˜๋ฆฌํ•ฉ๋‹ˆ๋‹ค.

  • ์ฑ„์šฉ ๊ณต๊ณ  ์ •๋ฆฌ(JOB_POSTING_INGEST)
  • ์ž์†Œ์„œ ๋ถ„์„(ANALYSIS)

์ž‘์—…์„ ํ์—์„œ ๊ตฌ๋…ํ•œ ๋’ค OpenAI ํ˜ธ์ถœ์„ ์ˆ˜ํ–‰ํ•˜๊ณ , ๊ฒฐ๊ณผ๋ฅผ Spring ๋‚ด๋ถ€ API๋กœ ์ €์žฅ/์™„๋ฃŒ ์ฒ˜๋ฆฌํ•˜๋Š” ๊ฒƒ์ด ์ด ๋ ˆํฌ์˜ ํ•ต์‹ฌ ์—ญํ• ์ž…๋‹ˆ๋‹ค.

1. ์—ญํ•  ์š”์•ฝ

๋ฉ”์ธ ๋ฐฑ์—”๋“œ(Spring)๋Š” ์‚ฌ์šฉ์ž ์š”์ฒญ์„ ๋ฐ›์€ ๋’ค ์ง์ ‘ ๊ธด AI ์ž‘์—…์„ ์ˆ˜ํ–‰ํ•˜์ง€ ์•Š์Šต๋‹ˆ๋‹ค. ๋Œ€์‹  RabbitMQ์— ์ž‘์—… ๋ฉ”์‹œ์ง€๋ฅผ ์ ์žฌํ•˜๊ณ , ์ด ์›Œ์ปค๊ฐ€ ๋ฉ”์‹œ์ง€๋ฅผ ์†Œ๋น„ํ•ฉ๋‹ˆ๋‹ค.

์›Œ์ปค๋Š” ๋‹ค์Œ ์ˆœ์„œ๋กœ ๋™์ž‘ํ•ฉ๋‹ˆ๋‹ค.

  1. RabbitMQ ํ๋ฅผ ๊ตฌ๋…ํ•ฉ๋‹ˆ๋‹ค.
  2. ๋ฉ”์‹œ์ง€๋ฅผ ์—ญ์ง๋ ฌํ™”ํ•ด ์ž‘์—… ํƒ€์ž…์„ ํŒ๋ณ„ํ•ฉ๋‹ˆ๋‹ค.
  3. Spring ๋‚ด๋ถ€ API์— ์ž‘์—… ์ƒํƒœ๋ฅผ RUNNING์œผ๋กœ ๋ฐ˜์˜ํ•ฉ๋‹ˆ๋‹ค.
  4. OpenAI๋ฅผ ํ˜ธ์ถœํ•ด ์ฑ„์šฉ ๊ณต๊ณ  ์ถ”์ถœ/๋ถ„๋ฅ˜/์ƒ์„ฑ ๋˜๋Š” ์ž์†Œ์„œ ๋ถ„์„์„ ์ˆ˜ํ–‰ํ•ฉ๋‹ˆ๋‹ค.
  5. ๊ฒฐ๊ณผ๋ฅผ Spring ๋‚ด๋ถ€ API์— ์ €์žฅํ•ฉ๋‹ˆ๋‹ค.
  6. ์ตœ์ข… ์™„๋ฃŒ ์ฝœ๋ฐฑ์„ ์ „๋‹ฌํ•ฉ๋‹ˆ๋‹ค.
  7. ์‹คํŒจ ์‹œ ์žฌ์‹œ๋„, DLQ ์ ์žฌ, recovery spool ๋ณต๊ตฌ๋ฅผ ์ˆ˜ํ–‰ํ•ฉ๋‹ˆ๋‹ค.

2. ์•„ํ‚คํ…์ฒ˜

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
Loading

3. ํ ์ ์žฌ / ๊ตฌ๋… / ๋ฐœํ–‰ ๋ฐฉ์‹

3.1 ๋ˆ„๊ฐ€ ํ์— ์ ์žฌํ•˜๋‚˜์š”?

ํ ์ ์žฌ๋Š” ์ด ๋ ˆํฌ๊ฐ€ ์•„๋‹ˆ๋ผ ๋ฉ”์ธ ๋ฐฑ์—”๋“œ(Spring)์—์„œ ์ˆ˜ํ–‰ํ•ฉ๋‹ˆ๋‹ค.

  • ์ฑ„์šฉ ๊ณต๊ณ  ์ž‘์—… APP_WORKER_JOB_POSTING_EXCHANGE + APP_WORKER_JOB_POSTING_ROUTING_KEY
  • ์ž์†Œ์„œ ๋ถ„์„ ์ž‘์—… APP_WORKER_ANALYSIS_EXCHANGE + APP_WORKER_ANALYSIS_ROUTING_KEY

์›Œ์ปค๋Š” ์ ์žฌ๋œ ๋ฉ”์‹œ์ง€๋ฅผ ์†Œ๋น„ํ•˜๋Š” ์ชฝ์ž…๋‹ˆ๋‹ค.

3.2 ์›Œ์ปค๋Š” ์–ด๋–ป๊ฒŒ ๊ตฌ๋…ํ•˜๋‚˜์š”?

์›Œ์ปค ์‹œ์ž‘ ์‹œ app/consumer.py ์˜ RabbitMqConsumer.start()๊ฐ€ ์‹คํ–‰๋˜๊ณ , ๋‚ด๋ถ€ async runtime์ด app/async_runtime.py ๋ฅผ ํ†ตํ•ด ๋‹ค์Œ ๋‘ ํ๋ฅผ ๋น„๋™๊ธฐ๋กœ ๊ตฌ๋…ํ•ฉ๋‹ˆ๋‹ค.

  • APP_WORKER_JOB_POSTING_QUEUE
  • APP_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_LIMIT
  • WORKER_ANALYSIS_CONCURRENCY_LIMIT
  • WORKER_JOB_POSTING_CONCURRENCY_LIMIT

๊ฐ ๋ฉ”์‹œ์ง€๋Š” ๊ณตํ†ต consumer๊ฐ€ task type๋ณ„ processor๋กœ dispatchํ•˜๊ณ , ์‹ค์ œ ๋™์‹œ ์ฒ˜๋ฆฌ ์Šฌ๋กฏ์€ app/concurrency.py ์˜ limiter๊ฐ€ ์ œ์–ดํ•ฉ๋‹ˆ๋‹ค. ์ด ๊ตฌ์กฐ ๋•๋ถ„์— analysis ์ ์ฒด๊ฐ€ jobposting ์ „์ฒด๋ฅผ ๋ง‰์ง€ ์•Š๋„๋ก ์ œํ•œ๊ฐ’์„ ๋ถ„๋ฆฌํ•ด ์šด์˜ํ•  ์ˆ˜ ์žˆ์Šต๋‹ˆ๋‹ค.

3.3 ์›Œ์ปค๊ฐ€ ๋ฐœํ–‰ํ•˜๋Š” ๋ฉ”์‹œ์ง€๋„ ์žˆ๋‚˜์š”?

์žˆ์Šต๋‹ˆ๋‹ค. ์ด ์›Œ์ปค๋Š” ์†Œ๋น„์ž์ด๋ฉด์„œ ์ผ๋ถ€ ์ƒํ™ฉ์—์„œ๋Š” ๋‹ค์‹œ RabbitMQ์— ๋ฉ”์‹œ์ง€๋ฅผ ๋ฐœํ–‰ํ•ฉ๋‹ˆ๋‹ค.

  • ์žฌ์‹œ๋„ ์‹œ ์›๋ณธ ๋ฉ”์‹œ์ง€์˜ retryCount๋ฅผ ์ฆ๊ฐ€์‹œ์ผœ ๊ฐ™์€ exchange/routing key๋กœ ์žฌ๋ฐœํ–‰
  • ์ตœ์ข… ์‹คํŒจ ์‹œ DLQ ํ๋กœ ๋ฐœํ–‰

์ฆ‰, ์ •์ƒ ์ฒ˜๋ฆฌ ๊ฒฐ๊ณผ๋Š” RabbitMQ๋กœ ๋‹ค์‹œ ๋ณด๋‚ด์ง€ ์•Š๊ณ  Spring ๋‚ด๋ถ€ API๋กœ ์ „๋‹ฌํ•˜๊ณ , ํ ์žฌ๋ฐœํ–‰์€ ์žฌ์‹œ๋„/์‹คํŒจ ์ฒ˜๋ฆฌ์—๋งŒ ์‚ฌ์šฉํ•ฉ๋‹ˆ๋‹ค.

4. ์ž‘์—… ํƒ€์ž…๋ณ„ ์ฒ˜๋ฆฌ ํ๋ฆ„

4.1 ์ฑ„์šฉ ๊ณต๊ณ  ์ •๋ฆฌ ์ž‘์—…

taskType=JOB_POSTING_INGEST

  1. Spring์ด ์ž‘์—… ๋ฉ”์‹œ์ง€๋ฅผ ํ์— ์ ์žฌํ•ฉ๋‹ˆ๋‹ค.
  2. ์›Œ์ปค๊ฐ€ ๋ฉ”์‹œ์ง€๋ฅผ ์†Œ๋น„ํ•ฉ๋‹ˆ๋‹ค.
  3. /api/internal/worker/job-postings/tasks/{taskId}/running ์œผ๋กœ ์ƒํƒœ๋ฅผ RUNNING ์ฒ˜๋ฆฌํ•ฉ๋‹ˆ๋‹ค.
  4. /api/internal/worker/job-postings/ingest/context ๋กœ ์ด๋ฏธ์ง€ URL ๋“ฑ ์ปจํ…์ŠคํŠธ๋ฅผ ์กฐํšŒํ•ฉ๋‹ˆ๋‹ค.
  5. OpenAI๋กœ ๊ณต๊ณ  ์ถ”์ถœ์„ ์ˆ˜ํ–‰ํ•ฉ๋‹ˆ๋‹ค.
  6. /api/internal/worker/job-postings/classification/candidates ๋กœ ๋ถ„๋ฅ˜ ํ›„๋ณด๋ฅผ ์กฐํšŒํ•ฉ๋‹ˆ๋‹ค.
  7. OpenAI๋กœ ํ›„๋ณด ์ค‘ ์†Œ๋ถ„๋ฅ˜๋ฅผ ์„ ํƒํ•ฉ๋‹ˆ๋‹ค.
  8. ์‹ ๋ขฐ๋„๊ฐ€ ๋‚ฎ์œผ๋ฉด ์ €์žฅ ์—†์ด /complete ๋กœ ์ข…๋ฃŒํ•ฉ๋‹ˆ๋‹ค.
  9. ์‹ ๋ขฐ๋„๊ฐ€ ์ถฉ๋ถ„ํ•˜๋ฉด OpenAI๋กœ ์ €์žฅ์šฉ ๊ณต๊ณ  ๋‚ด์šฉ์„ ์ƒ์„ฑํ•ฉ๋‹ˆ๋‹ค.
  10. /result ๋กœ ์ค‘๊ฐ„ ๊ฒฐ๊ณผ๋ฅผ ์ €์žฅํ•ฉ๋‹ˆ๋‹ค.
  11. /ingest/finalize ๋กœ ์ตœ์ข… ์ €์žฅ/์™„๋ฃŒ๋ฅผ ์š”์ฒญํ•ฉ๋‹ˆ๋‹ค.

4.2 ์ž์†Œ์„œ ๋ถ„์„ ์ž‘์—…

taskType=ANALYSIS

  1. Spring์ด ๋ถ„์„ ์ž‘์—… ๋ฉ”์‹œ์ง€๋ฅผ ํ์— ์ ์žฌํ•ฉ๋‹ˆ๋‹ค.
  2. ์›Œ์ปค๊ฐ€ ๋ฉ”์‹œ์ง€๋ฅผ ์†Œ๋น„ํ•ฉ๋‹ˆ๋‹ค.
  3. ํ ๋Œ€๊ธฐ ์‹œ๊ฐ„์ด APP_WORKER_ANALYSIS_QUEUE_TIMEOUT_MILLIS ๋ฅผ ์ดˆ๊ณผํ–ˆ๋Š”์ง€ ๋จผ์ € ๊ฒ€์‚ฌํ•ฉ๋‹ˆ๋‹ค.
  4. /api/internal/worker/analysis/tasks/{taskId}/running ์œผ๋กœ ์ƒํƒœ๋ฅผ RUNNING ์ฒ˜๋ฆฌํ•ฉ๋‹ˆ๋‹ค.
  5. /api/internal/worker/analysis/context ๋กœ ๋ถ„์„์— ํ•„์š”ํ•œ ๊ณต๊ณ /๋ฌธํ•ญ/๋‹ต๋ณ€ ๋ฐ์ดํ„ฐ๋ฅผ ์กฐํšŒํ•ฉ๋‹ˆ๋‹ค.
  6. OpenAI๋กœ ๋ถ„์„์„ ์ˆ˜ํ–‰ํ•ฉ๋‹ˆ๋‹ค.
  7. /result ๋กœ ๋ถ„์„ ๊ฒฐ๊ณผ๋ฅผ ์ €์žฅํ•ฉ๋‹ˆ๋‹ค.
  8. /complete ๋กœ ์ตœ์ข… ์™„๋ฃŒ ์ฒ˜๋ฆฌํ•ฉ๋‹ˆ๋‹ค.

5. ๋ฉ”์‹œ์ง€ ํ˜•์‹

5.1 ์ฑ„์šฉ ๊ณต๊ณ  ๋ฉ”์‹œ์ง€ ์˜ˆ์‹œ

{
  "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"
}

5.2 ์ž์†Œ์„œ ๋ถ„์„ ๋ฉ”์‹œ์ง€ ์˜ˆ์‹œ

{
  "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 ์— ์ •์˜๋˜์–ด ์žˆ์Šต๋‹ˆ๋‹ค.

6. ์žฌ์‹œ๋„, DLQ, ๋ณต๊ตฌ ์ „๋žต

์ด ์›Œ์ปค๋Š” ๋‹จ์ˆœ consume ํ›„ ์ข…๋ฃŒํ•˜์ง€ ์•Š๊ณ , ์‹คํŒจ ๋‚ด์„ฑ์„ ๊ฐ–์ถ˜ ๊ตฌ์กฐ๋กœ ์ž‘์„ฑ๋˜์–ด ์žˆ์Šต๋‹ˆ๋‹ค.

6.1 ์žฌ์‹œ๋„

  • RetryableWorkerError ๋ฐœ์ƒ ์‹œ ์žฌ์‹œ๋„ ๋Œ€์ƒ์œผ๋กœ ๊ฐ„์ฃผํ•ฉ๋‹ˆ๋‹ค.
  • retryCount + 1 ๊ฐ’์„ ๋‹ด์•„ ๊ฐ™์€ ํ ํ† ํด๋กœ์ง€๋กœ ์žฌ๋ฐœํ–‰ํ•ฉ๋‹ˆ๋‹ค.
  • ์ตœ๋Œ€ ์žฌ์‹œ๋„ ํšŸ์ˆ˜๋Š” ์ž‘์—…๋ณ„๋กœ ๋‹ค๋ฆ…๋‹ˆ๋‹ค.
    • ์ฑ„์šฉ ๊ณต๊ณ : WORKER_MAX_RETRY_COUNT
    • ์ž์†Œ์„œ ๋ถ„์„: APP_WORKER_ANALYSIS_MAX_RETRY_COUNT

6.2 DLQ

์•„๋ž˜ ๊ฒฝ์šฐ์—๋Š” DLQ๋กœ ๋ณด๋ƒ…๋‹ˆ๋‹ค.

  • ๋น„์žฌ์‹œ๋„ ์˜ค๋ฅ˜(NonRetryableWorkerError)
  • ์žฌ์‹œ๋„ ํšŸ์ˆ˜ ์ดˆ๊ณผ

์„ค์ • ํ‚ค๋Š” ๋‹ค์Œ๊ณผ ๊ฐ™์Šต๋‹ˆ๋‹ค.

  • APP_WORKER_JOB_POSTING_DLQ
  • APP_WORKER_ANALYSIS_DLQ

6.3 Recovery Spool

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 ํŒŒ์ผ์„ ๋‹ค์‹œ ์ฝ์–ด ๋ฏธ์ „๋‹ฌ ์™„๋ฃŒ ์š”์ฒญ์„ ์žฌ์ „์†กํ•ฉ๋‹ˆ๋‹ค.

๊ด€๋ จ ๊ตฌํ˜„์€ ๋‹ค์Œ ํŒŒ์ผ์— ์žˆ์Šต๋‹ˆ๋‹ค.

7. ๋‚ด๋ถ€ API ์—ฐ๋™ ํฌ์ธํŠธ

์ด ์›Œ์ปค๋Š” Spring ๊ณต๊ฐœ API๊ฐ€ ์•„๋‹ˆ๋ผ ๋‚ด๋ถ€ worker API๋ฅผ ํ˜ธ์ถœํ•ฉ๋‹ˆ๋‹ค. ๊ณตํ†ต ํ—ค๋”๋กœ X-Internal-Api-Key๋ฅผ ์‚ฌ์šฉํ•ฉ๋‹ˆ๋‹ค.

์ฃผ์š” ์—ฐ๋™ ์—”๋“œํฌ์ธํŠธ๋Š” ๋‹ค์Œ๊ณผ ๊ฐ™์Šต๋‹ˆ๋‹ค.

7.1 ์ฑ„์šฉ ๊ณต๊ณ  ์ž‘์—…

  • POST /api/internal/worker/job-postings/tasks/{taskId}/running
  • POST /api/internal/worker/job-postings/ingest/context
  • POST /api/internal/worker/job-postings/classification/candidates
  • POST /api/internal/worker/job-postings/tasks/{taskId}/result
  • POST /api/internal/worker/job-postings/ingest/finalize
  • POST /api/internal/worker/job-postings/tasks/{taskId}/retry
  • POST /api/internal/worker/job-postings/tasks/{taskId}/failed
  • GET /api/internal/worker/job-postings/tasks/{taskId}

7.2 ์ž์†Œ์„œ ๋ถ„์„ ์ž‘์—…

  • POST /api/internal/worker/analysis/tasks/{taskId}/running
  • POST /api/internal/worker/analysis/context
  • POST /api/internal/worker/analysis/tasks/{taskId}/result
  • POST /api/internal/worker/analysis/tasks/{taskId}/complete
  • POST /api/internal/worker/analysis/tasks/{taskId}/retry
  • POST /api/internal/worker/analysis/tasks/{taskId}/failed
  • GET /api/internal/worker/analysis/tasks/{taskId}

๊ตฌํ˜„์€ app/api_client.py ์— ์žˆ์œผ๋ฉฐ, ํ˜„์žฌ๋Š” sync ๋ฉ”์„œ๋“œ์™€ httpx.AsyncClient ๊ธฐ๋ฐ˜ async ๋ฉ”์„œ๋“œ๋ฅผ ํ•จ๊ป˜ ์ œ๊ณตํ•ฉ๋‹ˆ๋‹ค.

8. OpenAI ์ฒ˜๋ฆฌ ๋ฐฉ์‹

8.1 ์ฑ„์šฉ ๊ณต๊ณ  ์ž‘์—…

app/openai_client.py ์˜ JobPostingOpenAiWorker ๊ฐ€ ์•„๋ž˜ 3๋‹จ๊ณ„๋ฅผ ๋‹ด๋‹นํ•ฉ๋‹ˆ๋‹ค.

  1. extract ๊ณต๊ณ  ํ…์ŠคํŠธ ๋˜๋Š” ์ด๋ฏธ์ง€์—์„œ ๊ตฌ์กฐํ™” ์ •๋ณด ์ถ”์ถœ
  2. classify Spring์ด ์ค€ ๋ถ„๋ฅ˜ ํ›„๋ณด ์ค‘ ๊ฐ€์žฅ ์ ํ•ฉํ•œ ์†Œ๋ถ„๋ฅ˜ ์„ ํƒ
  3. generate ์ €์žฅ ๊ฐ€๋Šฅํ•œ ํ˜•ํƒœ์˜ ์ •์ œ๋œ ๊ณต๊ณ  ๋‚ด์šฉ ์ƒ์„ฑ

8.2 ์ž์†Œ์„œ ๋ถ„์„ ์ž‘์—…

AnalysisOpenAiWorker ๊ฐ€ ๋ฌธํ•ญ๋ณ„ ๋‹ต๋ณ€๊ณผ ๊ณต๊ณ  ๋งฅ๋ฝ์„ ๋ฐ”ํƒ•์œผ๋กœ ์ข…ํ•ฉ ์ ์ˆ˜์™€ ํ”ผ๋“œ๋ฐฑ์„ ์ƒ์„ฑํ•ฉ๋‹ˆ๋‹ค.

ํ˜„์žฌ OpenAI ํ˜ธ์ถœ๋„ sync/async wrapper๋ฅผ ๋ชจ๋‘ ๊ฐ€์ง€๊ณ  ์žˆ์œผ๋ฉฐ, ์‹ค์ œ consumer ๊ฒฝ๋กœ๋Š” async client๋ฅผ ์‚ฌ์šฉํ•ฉ๋‹ˆ๋‹ค.

๊ธฐ๋ณธ ๋ชจ๋ธ์€ ์•„๋ž˜ ํ™˜๊ฒฝ๋ณ€์ˆ˜๋กœ ์ œ์–ด๋ฉ๋‹ˆ๋‹ค.

  • OPENAI_JOB_POSTING_MODEL=gpt-4o-mini
  • OPENAI_ANALYSIS_MODEL=gpt-4.1-mini

9. ํ”„๋กœ์ ํŠธ ๊ตฌ์กฐ

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

10. ์‹คํ–‰ ๋ฐฉ๋ฒ•

10.0 ์š”๊ตฌ ์‚ฌํ•ญ

  • Python 3.12 ์ด์ƒ ๊ถŒ์žฅ
  • RabbitMQ ์ ‘๊ทผ ๊ฐ€๋Šฅ
  • Spring ๋‚ด๋ถ€ worker API ์ ‘๊ทผ ๊ฐ€๋Šฅ
  • OpenAI API ํ‚ค

์ด ๋ ˆํฌ๋Š” Dockerfile ๊ธฐ์ค€์œผ๋กœ Python 3.12 ํ™˜๊ฒฝ์—์„œ ์‹คํ–‰๋ฉ๋‹ˆ๋‹ค.

10.1 ๋กœ์ปฌ ์„ค์น˜

python3 -m venv .venv
source .venv/bin/activate
pip install -r requirements.txt

10.2 ํ•„์ˆ˜ ํ™˜๊ฒฝ๋ณ€์ˆ˜

์ตœ์†Œํ•œ ์•„๋ž˜ ๊ฐ’์€ ํ•„์š”ํ•ฉ๋‹ˆ๋‹ค.

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 ๋ฅผ ์ฐธ๊ณ ํ•˜๋ฉด ๋ฉ๋‹ˆ๋‹ค.

10.3 ์›Œ์ปค ์‹คํ–‰

uvicorn app.main:app --host 0.0.0.0 --port 8000
  • health endpoint GET /health
  • metrics endpoint GET /metrics

FastAPI ์•ฑ startup์—์„œ RabbitMQ consumer๋„ ํ•จ๊ป˜ ์‹œ์ž‘๋˜๋ฏ€๋กœ, ๋ณ„๋„ ๋ฐฑ๊ทธ๋ผ์šด๋“œ ์›Œ์ปค ํ”„๋กœ์„ธ์Šค๋ฅผ ์ถ”๊ฐ€๋กœ ๋„์šฐ์ง€ ์•Š์Šต๋‹ˆ๋‹ค.

11. Docker ๋ฐฐํฌ

์ด๋ฏธ์ง€๋Š” Dockerfile ๊ธฐ์ค€์œผ๋กœ ๋นŒ๋“œ๋˜๋ฉฐ ๊ธฐ๋ณธ ์‹คํ–‰ ๋ช…๋ น์€ ์•„๋ž˜์™€ ๊ฐ™์Šต๋‹ˆ๋‹ค.

CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"]

์šด์˜ ๋ฐฐํฌ ์˜ˆ์‹œ๋Š” deploy/docker-compose.worker.prod.yml ์— ์žˆ์Šต๋‹ˆ๋‹ค.

์ƒ์„ธ ์šด์˜ ๋ฌธ์„œ๋Š” ์•„๋ž˜ ํŒŒ์ผ์„ ์ฐธ๊ณ ํ•˜์„ธ์š”.

12. ๋กœ๊ทธ์™€ ๊ด€์ฐฐ ํฌ์ธํŠธ

๋กœ๊ทธ์—๋Š” ๊ธฐ๋ณธ์ ์œผ๋กœ ์•„๋ž˜ ์ปจํ…์ŠคํŠธ๊ฐ€ ํ•จ๊ป˜ ๋‚จ์Šต๋‹ˆ๋‹ค.

  • taskId
  • messageId
  • workerId
  • retryCount

์ •์ƒ ๊ธฐ๋™ ์‹œ ํ™•์ธํ•  ๋Œ€ํ‘œ ๋กœ๊ทธ:

  • Worker process started.
  • RabbitMQ consumer started.

13. ๋ฉ”ํŠธ๋ฆญ๊ณผ ํ•ด์„ ๊ธฐ์ค€

๊ธฐ์กด 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 ์ „ํ™˜๋Ÿ‰์ž…๋‹ˆ๋‹ค.

13.1 ์—๋Ÿฌ ๋ฐ fallback ์ •์ฑ…

  • 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 ์„ฑ๊ฒฉ ์˜ค๋ฅ˜๋Š” ๋น„์žฌ์‹œ๋„

13.2 Grafana / Prometheus ์ฟผ๋ฆฌ ์˜ˆ์‹œ

ํ ๋Œ€๊ธฐ 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])
  )
)

14. ํ…Œ์ŠคํŠธ

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 ๋กœ๊ทธ ํ•„๋“œ ์ •ํ•ฉ์„ฑ