[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으로 인해 엮여있는 객체들도 모두 초기화되는 모습이지요. 또한 별도 프로세스로 깨우는 등의 작업이 이루어집니다.
- 모델 실행기 초기화(i.e.,
ModelExecutor)- 태스크 큐/결과큐를 관리하는 멤버변수 초기화
- 워커 프로세스를 별도 프로세스로 구동. 이때 위의 멤버변수를 전달합니다. 그리고 동시에
ModelWorker.run()이 별도 프로세스에서 구동되죠
- 워크로드 매니저 초기화. 활성 요청을 트래킹하고 다음 요청배치를 결정
- 일반요청, 스트리밍요청에 따라 다른 큐를 사용
- 배치 사이즈 조절
- 요청 ID에 대해 처리할 수 있도록 하는 맵 객체(딕셔너리) 정의. 요청 ID별 조회, 맵 관리를 빠르게 구성하기 위함
- (CUDA가 세팅되어있다면) vLLM 객체 초기화.
- 스트리밍 요청을 처리하는 별도 루프 생성. 이 루프는 요청이 왔을 때 처리하는 쪽입니다. (아직 요청을 준 쪽이 없어 헷갈릴 수 있지만 잠시만 기다려주세요. 해당 로직은
workload_manager.add_request()/workload_manager.add_streaming_request()를 살펴보면 됩니다)- streaming 시퀀스를 가져와 워커 프로세스에 forward를 요청하고, 시퀀스 당 다음 토큰을 받아 클라이언트 큐로 분배
ModelExecutor() 초기화
모델 워커 프로세스와 태스크 큐, 결과 큐를 관리하는 객체를 초기화(아래의 setup_worker() 호출 참고)합니다.
개별 태스크 큐는 fork된 프로세스를 받을 수 있는 Queue() 로 구성됩니다.
WorkloadManager() 호출, 그리고 Sequence 클래스
Sequence클래스는 요청 ID와 프롬프트뿐 아니라 지금까지 생성한 출력, 완료 여부, 생성한 토큰 수가 저장됩니다. 스트리밍 요청이라면 클라이언트별asyncio.Queue와 FastAPI 이벤트 루프도 함께 참조합니다.WorkloadManager는 스케줄러입니다. 요청을Sequence단위로 등록하고, 다음 모델 실행에 포함할 배치를 고르며, 처리 중인 요청의 상태를 추적합니다.- 일반 요청과 스트리밍 요청은 각각 별도의 대기 큐(
incoming_queue,incoming_streaming_queue)와 활성 목록(active_sequences,active_streaming_sequences)으로 관리합니다. 대기 큐에서 꺼낸 요청은 활성 목록으로 이동하며, 실제 추론은ModelExecutor에 위임합니다. batch_size = 4는 한 번의 모델 실행 배치에 올려둘 활성Sequence의 최대 개수입니다. 요청이 4개보다 적으면 들어온 만큼만 처리합니다.sequence_map은 요청 ID를 키로Sequence를 바로 찾기 위한 맵입니다. 모델 결과를 원래 요청에 반영하거나 완료된 요청을 제거할 때 사용합니다.
- 일반 요청과 스트리밍 요청은 각각 별도의 대기 큐(
ModelExecutor.setup_worker(model_name) 구동
setup_worker()를 호출하여 워커 프로세스를multiprocessing패키지로 생성합니다. 이 시점에는 프로세스의 실행 방법만 준비합니다.start()로 자식 프로세스를fork하면 자식이ModelWorker.run()에 진입합니다. 실제facebook/opt-125m모델 초기화는 그 안에서ModelWorker(model_name)을 호출할 때 이루어집니다.- 아래와 같은 과정으로 진행되죠:
# 프로세스를 생성. # 이 프로세스는 ModelWorker의 `run()` 이라는 staticmethod 를 구동하도록 준비합니다 # args로 인자값을 준비합니다(태스크 큐와 결과 큐를 같이 담음) # daemon 옵션을 주어 부모 프로세스와 생명주기를 같이 합니다. (xNix의 daemonize가 아님에 유의) self.worker_process = _mp_context.Process( target=ModelWorker.run, args=(model_name, self.task_queue, self.result_queue), daemon=True, ) logger.debug("Starting worker process") self.worker_process.start() # 이때 모델 워커 프로세스를 fork합니다 logger.debug("Worker process started")
(참고) 인스턴스 없이 ModelWorker.run을 실행할 수 있는 이유
부모 프로세스에서 클래스 정의와 함수 정의를 하고, 자식 프로세스는 fork 를 떠서 갈 때 그 정의를 같이 가져갑니다.
실행시에는 Copy on write 방식에 따라, 자식 프로세스에서 호출할 때 메모리에서 복사하여 사용합니다.
model_worker모듈을 import할 때class ModelWorker:본문이 실행되어 클래스 객체와 그 안의 함수 정의가 먼저 만들어집니다.Process(target=ModelWorker.run, ...)은 자식 프로세스가 시작된 후 실행할 callable로 저장합니다.start()가 현재 프로세스를fork하면 자식은 이미 import된 모듈,ModelWorker클래스 정의,run함수 코드를 이어받습니다. 다만 이때도ModelWorker인스턴스가 생기지는 않습니다.- 이후
multiprocessing의 자식 프로세스 bootstrap이 다음과 같은 호출을 수행합니다.ModelWorker.run(model_name, task_queue, result_queue) run()은@staticmethod가 붙어 있으므로 인스턴스 없이 바로 호출할 수 있죠. 그래서 자식 프로세스 안에ModelWorker인스턴스를 만들고__init__()을 실행하여 모델을 로드합니다.worker = ModelWorker(model_name)
따라서 실행 순서는 다음과 같습니다.
부모: ModelWorker 클래스 정의를 import
→ 부모: Process(target=ModelWorker.run, ...) 구성
→ 부모: start()로 fork
→ 자식: ModelWorker.run(...) 실행
→ 자식: ModelWorker(model_name)으로 인스턴스 생성
→ 자식: __init__()에서 모델·토크나이저 로드
→ 자식: task_queue.get() 대기 루프 진입
마지막 단계의 task_queue.get() 대기가 실제로는 커널이 프로세스를 재우는 것임을 파이썬 표준 라이브러리, 시스템콜, 커널 소스, 구동 결과순으로 추적해 보았습니다.
이는 2주차 part 1 - 워커 프로세스는 왜 '서 있는가' (OS level)에 정리했습니다.
fork 결과 직접 확인하기
ModelWorker_id와run_function_id는 부모와 자식에서 동일합니다.
자식이 부모의 주소 공간을 Copy-on-Write 방식으로 계승했기 때문에 같은 가상 주소로 관찰됩니다.worker_id는 자식에서run()에 진입한 뒤 새로 생성된ModelWorker인스턴스의 주소입니다.
부모 프로세스에서 워커를 시작하기 직전의 로그입니다.
[fork-observe][parent]
pid=245967 ppid=3310 start_method=fork
ModelWorker_id=0x1d41aba0
run_function_id=0x7f1af79f37e0
fork 후 ModelWorker.run()에 진입한 자식 프로세스의 로그입니다.
Worker process started: parent_pid=245967 child_pid=246021
[fork-observe][child-before-instance]
pid=246021 ppid=245967 start_method=fork
ModelWorker_id=0x1d41aba0
run_function_id=0x7f1af79f37e0
자식 프로세스 안에서 worker = ModelWorker(model_name)을 실행한 뒤에는 별도의 인스턴스 ID가 확인됩니다.
[fork-observe][child-after-instance]
pid=246021
worker_id=0x7f1bf5b4b410
worker_class_id=0x1d41aba0
worker_class_is_ModelWorker=True
- 부모와 자식의 PID는 각각
245967,246021로 서로 다릅니다.
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()
Thread(...)는requests_processing_loop를 실행할 스레드 객체를 구성하고,start()가 실제 백그라운드 스레드를 시작합니다.requests_processing_loop()는while True로 구성되어 앱이 살아 있는 동안 계속 실행됩니다.daemon=True이므로 메인 프로세스가 종료될 때 이 스레드도 함께 종료됩니다. 현재 코드에는 스레드만을 위한 별도의 정상 종료 조건은 없습니다.- 루프는
WorkloadManager에서 처리할 스트리밍Sequence가 있는지 확인합니다.
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()에 서 있으면, 결과는 순전히 읽기 잠금을 먼저 잡는 순서로 배달됩니다.
전제가 깨지는 순간
깨지려면 세 조건이 동시에 맞아야 합니다.
- 두 스레드의 작업이
task_queue에 겹쳐 들어가 있고, - 결과가 도착하는 순간 그 결과의 주인이 아직
get()에 도착하지 못했으며, - 그 틈에 상대 스레드가 먼저 읽기 잠금을 잡아야 합니다.
실제 코드에서 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회 무발현이라는 결과 자체가, 이런 버그가 테스트를 통과하고 프로덕션까지 살아남는 이유를 보여줍니다.
sleep(1.0)으로 put과 get 사이의 틈을 인위적으로 벌려, 선점이 일어났을 때 벌어지는 일을 결정적으로 재현합니다.
import multiprocessing as mp
import threading
import time
ctx = mp.get_context("fork")
task_q = ctx.Queue()
result_q = ctx.Queue()
def worker(task_q, result_q):
"""ModelWorker.run() 흉내. 받은 순서대로 처리해서 결과를 넣는다."""
while True:
item = task_q.get()
if item is None:
break
batch, is_streaming = item
# 일반 배치(50토큰 생성)는 오래, forward 1회는 짧게 걸린다고 가정
time.sleep(0.5 if not is_streaming else 0.05)
result_q.put(("stream" if is_streaming else "complete", batch))
def execute_batch(batch):
"""ModelExecutor.execute_batch 흉내 (일반 요청)."""
task_q.put((batch, False))
# put과 get 사이에 이 스레드가 오래 선점당했다고 가정.
# 실제 코드에서는 이 틈이 훨씬 짧지만 0이 아니라는 게 요점.
time.sleep(1.0)
return result_q.get()
def execute_forward_batch(batch):
"""ModelExecutor.execute_forward_batch 흉내 (스트리밍 루프)."""
task_q.put((batch, True))
result_type, results = result_q.get()
if result_type != "stream":
raise RuntimeError(
f"UnexpectedResultType: {result_type!r} 를 받음 (내용: {results})"
)
return results
def generate_handler():
"""/generate 핸들러 역할 (요청 스레드)."""
res = execute_batch(["generate 배치"])
print(f"[generate 스레드] 받은 결과: {res}")
def streaming_loop():
"""requests_processing_loop 역할 (백그라운드 스레드)."""
time.sleep(0.1) # 일반 요청 직후 스트리밍 배치가 나가는 상황
try:
res = execute_forward_batch(["streaming 배치"])
print(f"[streaming 스레드] 받은 결과: {res}")
except RuntimeError as e:
print(f"[streaming 스레드] 예외 발생: {e}")
if __name__ == "__main__":
p = ctx.Process(target=worker, args=(task_q, result_q), daemon=True)
p.start()
t1 = threading.Thread(target=generate_handler)
t2 = threading.Thread(target=streaming_loop)
t1.start(); t2.start()
t1.join(); t2.join()
task_q.put(None)
실행 결과, 두 스레드가 서로의 결과를 맞바꿔 가집니다.
[streaming 스레드] 예외 발생: UnexpectedResultType: 'complete' 를 받음 (내용: ['generate 배치'])
[generate 스레드] 받은 결과: ('stream', ['streaming 배치'])
어떻게 고칠 것인가
execute_forward_batch()의 result_type != 'stream' 검사는 뒤바뀜을 "감지"할 뿐 "방지"하지 못하고, 그나마 일반 경로(execute_batch)에는 검사조차 없습니다. 수정 방향은 세 가지입니다.
- 경로별 큐 쌍 분리: 일반용과 스트리밍용 task/result 큐를 따로 둡니다. 호출자가 정확히 2개인 지금 구조에서는 간단하고 확실하지만, 워커가 큐 여러 개를 살펴야 하고 호출자가 늘 때마다 큐도 늘어납니다.
- 요청 ID 상관관계(correlation): 결과에 "누구의 결과인지"를 싣고 받는 쪽이 자기 것만 집게 합니다. 배치에 이미
request_id가 있으니 자연스러운 확장이고, vLLM을 포함한 실전 서빙 시스템이 쓰는 일반해입니다. - 왕복 전체를 락으로 직렬화:
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 객체 초기화를 합니다.