Ray Serve 비동기 추론
Ray Serve 비동기 추론
영상 처리나 대용량 문서 인덱싱처럼 HTTP 타임아웃보다 오래 걸리는 추론 요청은 그냥 기다리게 두면 API가 응답을 못 내죠. Ray Serve의 **비동기 추론(asynchronous inference)**은 이런 작업을 백그라운드 큐에 넣어 나중에 처리하고, 사용자에게는 즉시 빠른 응답을 돌려줘요. 요청 수명과 계산 시간을 분리하면서도 Serve의 스케일링을 그대로 활용할 수 있어요.
:::warning 경고 이 API는 alpha 단계라 안정 버전이 되기 전에 변경될 수 있어요. :::
왜 비동기 추론인가
Ray Serve 사용자가 오래 걸리는 API 요청을 비동기로 처리할 방법이 필요해요. 일부 추론 워크로드(영상 처리, 대용량 문서 인덱싱 등)는 일반적인 HTTP 타임아웃보다 오래 걸리기 때문에, 사용자가 요청을 제출하면 시스템은 작업을 백그라운드 큐에 넣고 즉시 빠른 응답을 반환해야 해요. 이렇게 하면 작업이 비동기로 실행되는 동안 요청 수명과 계산 시간이 분리되고, Serve의 스케일링도 그대로 활용할 수 있어요.
사용 사례
흔한 사용 사례는 영상 추론(긴 영상의 트랜스코딩·탐지·전사)과 대용량 파일이나 배치를 수집·파싱·벡터화하는 문서 인덱싱 파이프라인이에요. 더 넓게 보면 즉시 결과가 필요 없는 오래 걸리는 AI/ML 워크로드라면 비동기로 실행하는 게 좋아요.
핵심 개념
- @task_consumer — 큐에서 태스크를 소비·실행하는 Serve deployment예요. 태스크 프로세서를 구성하는
TaskProcessorConfig파라미터가 필요해요. 기본적으로 Celery 태스크 프로세서를 쓰지만 직접 구현을 제공할 수도 있어요. - @task_handler —
@task_consumer클래스 안의 메서드에 적용하는 데코레이터예요. 각 핸들러는name=...으로 처리할 태스크를 선언하고,name을 생략하면 메서드 함수명이 태스크 이름이 돼요. 소비자의 구성된 큐(TaskProcessorConfig로 설정)에서 그 이름을 가진 모든 태스크가 이 메서드로 라우팅돼 실행돼요.
구성 요소와 API
TaskProcessorConfig
큐 이름, 어댑터(기본은 Celery), 어댑터 설정, 재시도 한도, 데드-레터 큐를 포함한 태스크 프로세서를 구성해요.
from ray.serve.schema import TaskProcessorConfig, CeleryAdapterConfig
processor_config = TaskProcessorConfig(
queue_name="my_queue",
# Optional: Override default adapter string (default is Celery)
# adapter="ray.serve.task_processor.CeleryTaskProcessorAdapter",
adapter_config=CeleryAdapterConfig(
broker_url="redis://localhost:6379/0", # Or "filesystem://" for local testing
backend_url="redis://localhost:6379/1", # Result backend (optional for fire-and-forget)
),
max_retries=5,
failed_task_queue_name="failed_tasks", # Application errors after retries
)
:::note 참고
filesystem broker는 로컬 테스트 전용이고 기능이 제한적이에요. 예를 들어 cancel_tasks를 지원하지 않아요. 프로덕션 배포에서는 Redis나 RabbitMQ 같은 프로덕션급 broker를 쓰세요. 지원되는 broker 전체 목록은 Celery broker 문서를 참고하세요.
:::
@task_consumer
제공된 TaskProcessorConfig로 Serve deployment를 태스크 소비자로 바꾸는 데코레이터예요.
from ray import serve
from ray.serve.task_consumer import task_consumer
@serve.deployment
@task_consumer(task_processor_config=processor_config)
class SimpleConsumer:
pass
@task_handler
소비자에 메서드를 named task handler로 등록하는 데코레이터예요.
from ray.serve.task_consumer import task_handler, task_consumer
@serve.deployment
@task_consumer(task_processor_config=processor_config)
class SimpleConsumer:
@task_handler(name="process_request")
def process_request(self, data):
return f"processed: {data}"
:::note 참고
Ray Serve는 현재 동기 핸들러만 지원해요. async def 핸들러를 선언하면 NotImplementedError가 발생해요.
:::
instantiate_adapter_from_config
주어진 TaskProcessorConfig에 대해 태스크 프로세서 어댑터 인스턴스를 반환하는 팩토리 함수예요. 반환된 객체로 태스크를 enqueue하고, 상태를 조회하고, 메트릭을 가져오는 등의 작업을 할 수 있어요.
from ray.serve.task_consumer import instantiate_adapter_from_config
adapter = instantiate_adapter_from_config(task_processor_config=processor_config)
# Enqueue synchronously (returns TaskResult)
result = adapter.enqueue_task_sync(task_name="process_request", args=["hello"])
# Later, fetch status synchronously
status = adapter.get_task_status_sync(result.id)
:::note 참고
@serve.deployment 데코레이터에 지정된 모든 Ray actor 옵션(num_gpus, num_cpus, resources 등)이 태스크 소비자 레플리카에 적용돼요. 이 덕에 태스크 처리 워크로드에 특정 하드웨어 리소스를 할당할 수 있어요.
:::
End-to-end 예시: 문서 인덱싱
이 예시는 프로세서를 구성하고, 핸들러가 있는 소비자를 만들고, ingress deployment에서 태스크를 enqueue하고, 태스크 상태를 확인하는 전체 흐름을 보여줘요.
import io
import logging
import requests
from fastapi import FastAPI
from pydantic import BaseModel, HttpUrl
from PyPDF2 import PdfReader
from ray import serve
from ray.serve.schema import CeleryAdapterConfig, TaskProcessorConfig
from ray.serve.task_consumer import (
instantiate_adapter_from_config,
task_consumer,
task_handler,
)
logger = logging.getLogger("ray.serve")
fastapi_app = FastAPI(title="Async PDF Processing API")
TASK_PROCESSOR_CONFIG = TaskProcessorConfig(
queue_name="pdf_processing_queue",
adapter_config=CeleryAdapterConfig(
broker_url="redis://127.0.0.1:6379/0",
backend_url="redis://127.0.0.1:6379/0",
),
max_retries=3,
failed_task_queue_name="failed_pdfs",
)
class ProcessPDFRequest(BaseModel):
pdf_url: HttpUrl
max_summary_paragraphs: int = 3
@serve.deployment(num_replicas=2, max_ongoing_requests=5)
@task_consumer(task_processor_config=TASK_PROCESSOR_CONFIG)
class PDFProcessor:
"""Background worker that processes PDF documents asynchronously."""
@task_handler(name="process_pdf")
def process_pdf(self, pdf_url: str, max_summary_paragraphs: int = 3):
"""Download PDF, extract text, and generate summary."""
try:
response = requests.get(pdf_url, timeout=30)
response.raise_for_status()
pdf_reader = PdfReader(io.BytesIO(response.content))
if not pdf_reader.pages:
raise ValueError("PDF contains no pages")
full_text = "\n".join(
page.extract_text() for page in pdf_reader.pages if page.extract_text()
)
if not full_text.strip():
raise ValueError("PDF contains no extractable text")
paragraphs = [p.strip() for p in full_text.split("\n\n") if p.strip()]
summary = "\n\n".join(paragraphs[:max_summary_paragraphs])
return {
"status": "success",
"pdf_url": pdf_url,
"page_count": len(pdf_reader.pages),
"word_count": len(full_text.split()),
"summary": summary,
}
except requests.exceptions.RequestException as e:
raise ValueError(f"Failed to download PDF: {str(e)}")
except Exception as e:
raise ValueError(f"Failed to process PDF: {str(e)}")
@serve.deployment()
@serve.ingress(fastapi_app)
class AsyncPDFAPI:
"""HTTP API for submitting and checking PDF processing tasks."""
def __init__(self, task_processor_config: TaskProcessorConfig, handler):
self.adapter = instantiate_adapter_from_config(task_processor_config)
@fastapi_app.post("/process")
def process_pdf(self, request: ProcessPDFRequest):
"""Submit a PDF processing task and return task_id immediately."""
task_result = self.adapter.enqueue_task_sync(
task_name="process_pdf",
kwargs={
"pdf_url": str(request.pdf_url),
"max_summary_paragraphs": request.max_summary_paragraphs,
},
)
return {
"task_id": task_result.id,
"status": task_result.status,
"message": "PDF processing task submitted successfully",
}
@fastapi_app.get("/status/{task_id}")
def get_status(self, task_id: str):
"""Get task status and results."""
status = self.adapter.get_task_status_sync(task_id)
return {
"task_id": task_id,
"status": status.status,
"result": status.result if status.status == "SUCCESS" else None,
"error": str(status.result) if status.status == "FAILURE" else None,
}
app = AsyncPDFAPI.bind(TASK_PROCESSOR_CONFIG, PDFProcessor.bind())
이 예시에서:
DocumentIndexingConsumer는document_indexing_queue큐에서 태스크를 읽어 처리해요.API는enqueue_task_sync로 태스크를 enqueue하고get_task_status_sync로 상태를 조회해요.consumer를API.__init__에 넘기는 것으로 두 deployment가 Serve application 그래프의 일부가 돼요.
동시성과 신뢰성
소비자 deployment의 max_ongoing_requests로 동시성을 관리해요. 이 값은 각 레플리카가 동시에 처리할 수 있는 태스크 수를 제한해요. at-least-once 전달을 위해 어댑터는 핸들러가 성공적으로 완료된 후에만 태스크를 ack 해야 해요. 실패한 태스크는 max_retries까지 재시도되고, 다 소진하면 설정된 경우 failed-task DLQ로 라우팅돼요. 기본 Celery 어댑터는 성공 시 ack해서 at-least-once 처리를 제공해요.
데드 레터 큐 (DLQ)
데드 레터 큐는 두 종류의 문제 태스크를 처리해요:
- 처리 불가 태스크: 일치하는 핸들러가 없는 태스크는
unprocessable_task_queue_name이 설정돼 있으면 그 큐로 라우팅돼요. - 실패 태스크: 재시도 소진 후 애플리케이션 예외를 일으키거나, 인자가 일치하지 않거나, 기타 오류가 있는 태스크는
failed_task_queue_name이 설정돼 있으면 그 큐로 라우팅돼요.
제약 사항
- Ray Serve는 동기
@task_handler메서드만 지원해요. - 외부(비-Serve) 워커는 범위 밖이에요. 모든 소비자는 Serve deployment로 실행돼요.
- 전달 보장은 설정된 broker에 달려 있어요. 결과 백엔드를 구성하지 않으면 결과는 선택 사항이에요.
:::note 참고
이 가이드의 API는 ray.serve.schema와 ray.serve.task_consumer의 alpha 인터페이스를 반영한 것이에요.
:::