[CloudNeta] Hands-On LLM Serving 2주차 part 1 - 코드 분석 (1)

이 문서는 2주차 part 1 - 모델 서빙 시스템 설계에서 다룬 단일 모델 서비스의 코드 분석 편입니다.

코드의 전반적인 구성은 어렵지 않습니다. API 서버가 구성되어있고, 구동 시 LLMEngine() 을 실행하며 아래 요소를 초기화합니다.

서버 구동

서버 구동은 uv 로 진행합니다.

uv run uvicorn main:app --host 0.0.0.0 --port 8000

초기화과정

LLMEngine() 객체를 초기화하여 전역적으로 사용합니다. 그리고 코드구조와 대비하기 위해 vLLM 객체도 초기화합니다.

이는 직접 구동되는 최소한의 예시와 vLLM 구성이 이를 어떻게 간략화하는지를 대비시켜주는 정도로 보시면 좋을 것 같습니다.

LLMEngine() 을 호출하면

아래 과정이 일어납니다. 앞선 그림의 LLMEngine으로 인해 엮여있는 객체들도 모두 초기화되는 모습이지요. 또한 별도 프로세스로 깨우는 등의 작업이 이루어집니다.

  1. 모델 실행기 초기화(i.e., ModelExecutor)
    • 태스크 큐/결과큐를 관리하는 멤버변수 초기화
    • 워커 프로세스를 별도 프로세스로 구동. 이때 위의 멤버변수를 전달합니다. 그리고 동시에 ModelWorker.run() 이 별도 프로세스에서 구동되죠
  2. 워크로드 매니저 초기화. 활성 요청을 트래킹하고 다음 요청배치를 결정
    • 일반요청, 스트리밍요청에 따라 다른 큐를 사용
    • 배치 사이즈 조절
    • 요청 ID에 대해 처리할 수 있도록 하는 맵 객체(딕셔너리) 정의. 요청 ID별 조회, 맵 관리를 빠르게 구성하기 위함
  3. (CUDA가 세팅되어있다면) vLLM 객체 초기화.
  4. 스트리밍 요청을 처리하는 별도 루프 생성. 이 루프는 요청이 왔을 때 처리하는 쪽입니다. (아직 요청을 준 쪽이 없어 헷갈릴 수 있지만 잠시만 기다려주세요. 해당 로직은 workload_manager.add_request() / workload_manager.add_streaming_request() 를 살펴보면 됩니다)
    • streaming 시퀀스를 가져와 워커 프로세스에 forward를 요청하고, 시퀀스 당 다음 토큰을 받아 클라이언트 큐로 분배

ModelExecutor() 초기화

모델 워커 프로세스와 태스크 큐, 결과 큐를 관리하는 객체를 초기화(아래의 setup_worker() 호출 참고)합니다.

개별 태스크 큐는 fork된 프로세스를 받을 수 있는 Queue() 로 구성됩니다.

WorkloadManager() 호출, 그리고 Sequence 클래스

ModelExecutor.setup_worker(model_name) 구동

(참고) 인스턴스 없이 ModelWorker.run을 실행할 수 있는 이유

핵심 정리

부모 프로세스에서 클래스 정의와 함수 정의를 하고, 자식 프로세스는 fork 를 떠서 갈 때 그 정의를 같이 가져갑니다.
실행시에는 Copy on write 방식에 따라, 자식 프로세스에서 호출할 때 메모리에서 복사하여 사용합니다.

따라서 실행 순서는 다음과 같습니다.

부모: ModelWorker 클래스 정의를 import
  → 부모: Process(target=ModelWorker.run, ...) 구성
  → 부모: start()로 fork
  → 자식: ModelWorker.run(...) 실행
  → 자식: ModelWorker(model_name)으로 인스턴스 생성
  → 자식: __init__()에서 모델·토크나이저 로드
  → 자식: task_queue.get() 대기 루프 진입
워커는 왜 CPU를 안 쓰고 "서 있는가"

마지막 단계의 task_queue.get() 대기가 실제로는 커널이 프로세스를 재우는 것임을 파이썬 표준 라이브러리, 시스템콜, 커널 소스, 구동 결과순으로 추적해 보았습니다.
이는 2주차 part 1 - 워커 프로세스는 왜 '서 있는가' (OS level)에 정리했습니다.

requests_processing_loop() 를 스레드로 구동

스트리밍 요청을 처리하는 루프는 FastAPI 이벤트 루프와 분리된 백그라운드 스레드에서 실행됩니다. 이 스레드는 요청을 배치로 구성하고 모델 워커 프로세스와 통신하지만, SSE 응답을 직접 전송하지는 않습니다.

1. 객체 생성

# Start processing loop in a separate thread
self.thread = threading.Thread(
    target=self.requests_processing_loop,
    daemon=True,
)
self.thread.start()
active_sequences = self.workload_manager.get_next_batch(is_streaming=True)
if not active_sequences:
    time.sleep(0.1)
    continue

요청이 없다면 0.1초 쉬었다가 다시 확인합니다. 요청이 있다면 incoming_streaming_queue의 시퀀스를 active_streaming_sequences로 옮긴 뒤 모델에 보낼 배치를 만듭니다.

2. execute_forward_batch() 이후

활성 시퀀스의 현재 프롬프트와 요청 ID를 모델 입력 형태로 만든 뒤 ModelExecutor에 전달합니다.

prompts = [
    {'prompt': seq.prompt, 'request_id': seq.id}
    for seq in active_sequences
]
prompts_results = self.model_executor.execute_forward_batch(prompts)

execute_forward_batch()는 배치와 스트리밍 여부를 multiprocessing.Queue에 넣습니다.

self.task_queue.put((prompts, True))
result_type, results = self.result_queue.get()

이때 자식 프로세스의 ModelWorker.run()은 이미 아래 코드에서 작업을 기다리고 있습니다.

batch_data = task_queue.get()

task_queue.put()이 호출되면 get()에서 대기하던 자식 프로세스가 배치를 받아 worker.generate_forward_batch(batch)를 실행합니다. 한 번의 호출에서 활성 시퀀스마다 다음 토큰 하나를 생성하고, 결과를 result_queue에 넣습니다.

백그라운드 스레드                         모델 워커 프로세스

task_queue.put((prompts, True))
                 ───────────────────────▶ task_queue.get()
                                           generate_forward_batch()
result_queue.get()
                 ◀─────────────────────── result_queue.put(...)

put()이 자식 프로세스로 제어권을 넘기는 것은 아닙니다. 부모의 백그라운드 스레드와 자식 프로세스는 독립적으로 실행되고, result_queue.get()에서 백그라운드 스레드만 결과가 올 때까지 기다립니다. FastAPI 이벤트 루프는 별도로 계속 실행됩니다.

3. lockstep: result_queue.get()은 어떻게 "내 결과"인 걸 아는가

위 코드를 보면 이상한 점이 하나 있습니다. result_queue.get()이 가져온 결과가 방금 내가 put()한 배치의 결과라는 보장은 어디에 있을까요? 결과에는 ('complete', ...) / ('stream', ...) 태그만 있지, "어느 호출자의 결과인지"를 대조하는 로직이 없습니다.

답은 매칭 로직이 없어도 되는 구조라는 것입니다. put 하나에 get 하나, 미결(in-flight) 작업이 항상 1개뿐인 1:1 lockstep이라서, "지금 큐에 들어올 결과는 방금 내가 넣은 작업의 결과일 수밖에 없다"는 전제 위에 서 있습니다.

그런데 이 전제는 계층을 나눠 보면 생각보다 약합니다.

계층 분리 여부
대기 큐 (incoming_queue / incoming_streaming_queue) 분리됨
배치 구성 (get_next_batch()의 is_streaming 분기) 분리됨
워커와의 통신 (task_queue / result_queue) 공유. 한 쌍뿐

요청을 "대기시키는" 계층은 일반/스트리밍이 분리되어 있지만, 워커와 "통신하는" 계층은 한 쌍의 큐를 같이 씁니다. execute_batch()(일반 요청, 요청 핸들러에서 실행)와 execute_forward_batch()(스트리밍, 백그라운드 스레드에서 실행)가 같은 task_queue에 넣고 같은 result_queue에서 꺼냅니다. 두 스레드가 동시에 result_queue.get()에 서 있으면, 결과는 순전히 읽기 잠금을 먼저 잡는 순서로 배달됩니다.

전제가 깨지는 순간

깨지려면 세 조건이 동시에 맞아야 합니다.

  1. 두 스레드의 작업이 task_queue에 겹쳐 들어가 있고,
  2. 결과가 도착하는 순간 그 결과의 주인이 아직 get()에 도착하지 못했으며,
  3. 그 틈에 상대 스레드가 먼저 읽기 잠금을 잡아야 합니다.

실제 코드에서 put()과 get() 사이는 바이트코드 몇 개, 마이크로초 단위입니다. GIL 스위치 주기(기본 5ms)를 생각하면 이 틈에 선점당할 확률이 매우 낮아서, 평상시에는 사실상 터지지 않습니다. 하지만 이것은 보장이 아니라 타이밍 요행입니다.

결과: 무부하에서는 멀쩡, CPU 포화에서는 28% 재현

스트리밍이 진행되는 동안 /basic_generate를 반복해서 끼워 넣는 실험을 했습니다.

조건 결과
유휴 시스템, 40라운드 0회 발생
CPU 16코어를 busy loop로 포화, 50라운드 14회 발생 (28%)

터졌을 때 서버 로그에는 두 피해자가 정확히 서로의 결과를 쥔 모습이 찍힙니다.

스트리밍 스레드는 일반 요청의 결과('complete')를 집어가서 예외가 나고:

ERROR - Error in processing loop
Traceback (most recent call last):
  File ".../llm/llm.py", line 62, in requests_processing_loop
    prompts_results = self.model_executor.execute_forward_batch(prompts)
  File ".../llm/model_executor.py", line 77, in execute_forward_batch
    raise UnexpectedResultTypeError("Unexpected result type from worker")
llm.exceptions.UnexpectedResultTypeError: Unexpected result type from worker

/basic_generate 핸들러는 스트리밍 결과를 받아 generated_text 키를 찾다가 500을 반환합니다.

DEBUG - Received results from worker: ('stream', [{'request_id': 'd38b00ca-...', 'token': ' I', 'is_finished': False}])
INFO:  127.0.0.1:55062 - "POST /basic_generate HTTP/1.1" 500 Internal Server Error
...
KeyError: 'generated_text'

"부하가 없을 땐 멀쩡하다가 부하가 걸리면 확률적으로 터지는" 전형적인 race condition입니다. 무부하 40회 무발현이라는 결과 자체가, 이런 버그가 테스트를 통과하고 프로덕션까지 살아남는 이유를 보여줍니다.

어떻게 고칠 것인가

execute_forward_batch()의 result_type != 'stream' 검사는 뒤바뀜을 "감지"할 뿐 "방지"하지 못하고, 그나마 일반 경로(execute_batch)에는 검사조차 없습니다. 수정 방향은 세 가지입니다.

  1. 경로별 큐 쌍 분리: 일반용과 스트리밍용 task/result 큐를 따로 둡니다. 호출자가 정확히 2개인 지금 구조에서는 간단하고 확실하지만, 워커가 큐 여러 개를 살펴야 하고 호출자가 늘 때마다 큐도 늘어납니다.
  2. 요청 ID 상관관계(correlation): 결과에 "누구의 결과인지"를 싣고 받는 쪽이 자기 것만 집게 합니다. 배치에 이미 request_id가 있으니 자연스러운 확장이고, vLLM을 포함한 실전 서빙 시스템이 쓰는 일반해입니다.
  3. 왕복 전체를 락으로 직렬화: put()부터 get()까지를 뮤텍스로 묶어 한 번에 한 왕복만 허용합니다. 수정은 가장 짧지만 두 경로가 서로를 기다리게 되어 처리량을 희생합니다.

4. client_stream과 loop

client_stream과 loop는 /generate_stream 요청의 event_generator()가 시작될 때 생성·등록됩니다.

queue = asyncio.Queue()
seq_id = self.workload_manager.add_streaming_request(prompt, queue, loop)
client_stream의 실제 정체

client_stream은 네트워크 스트림 자체가 아니라, 해당 요청으로 전달할 토큰을 담는 요청별 asyncio.Queue입니다. loop는 그 HTTP 요청을 처리하는 FastAPI/Uvicorn 이벤트 루프입니다.

requests_processing_loop()는 일반 스레드에서 실행되므로 await seq.client_stream.put(...)을 직접 호출할 수 없습니다. 따라서 run_coroutine_threadsafe()를 사용해 기존 FastAPI 이벤트 루프에 큐 삽입 작업을 예약합니다.

asyncio.run_coroutine_threadsafe(
    seq.client_stream.put(
        json.dumps({
            "token": result['token'],
            "sequence_id": result['request_id'],
        })
    ),
    seq.loop,
)

별도의 이벤트 루프를 새로 만드는 것이 아니라, Sequence에 저장해 둔 기존 이벤트 루프에서 asyncio.Queue.put() 코루틴을 실행하는 방식입니다. 큐에 데이터가 들어오면 event_generator()의 await queue.get()이 깨어나고, yield한 문자열이 StreamingResponse를 통해 SSE 본문으로 전송됩니다.

data = await queue.get()
if data is None:
    break
yield f"data: {data}\n\n"

생성이 끝나면 토큰 대신 None을 큐에 넣습니다. event_generator()는 이를 종료 신호로 인식하고 스트리밍을 마친 뒤 해당 Sequence를 정리합니다.

vLLM 객체 초기화

CUDA 사용이 가능하다면 VLLM 객체 초기화를 합니다.