[CloudNeta] Hands-On LLM Serving 2주차 - 추적: request_id의 여정 (스트리밍 경로)
이 문서는 2주차 part 1 - 모델 서빙 시스템 설계의 /generate_stream 경로를 코드 순서대로 따라간 추적 기록입니다. 일괄 경로는 generate의 대기 루프에서 다룹니다.
스트리밍 경로가 복잡해 보이는 이유는 단계가 많아서가 아닙니다. 데이터가 실행 경계를 네 번 넘기 때문입니다. 넘을 때마다 종류가 다른 큐가 필요하고, 돌아올 때 누구 것인지 알아야 해서 request_id가 끝까지 따라다닙니다.
네 개의 실행 구역
| 구역 | 무엇이 도는가 | 이 구역에만 있는 것 |
|---|---|---|
| Client | 브라우저 · curl | HTTP 연결 |
| 이벤트 루프 | 코루틴 (event_generator) |
asyncio.Queue |
| 워커 스레드 | requests_processing_loop 데몬 |
배치 슬롯 |
| 워커 프로세스 | ModelWorker · torch |
토큰 ID |
HTTP 경계 바깥은 처음부터 끝까지 사람이 읽는 문자열입니다. 토큰 ID([2, 133, 2119] 같은 숫자)는 맨 아래 구역을 절대 벗어나지 않고, 위로 올라올 때는 이미 문자열로 풀린 뒤입니다.
한 바퀴 전체
아래는 단 한 번의 forward에 해당하는 왕복입니다. 토큰 하나를 만들려고 이만큼을 돌고, 이걸 21번 반복합니다.
큐가 세 종류인 이유
코드에 queue라는 이름이 여러 번 나오는데 전부 다른 물건입니다. 무엇과 무엇을 잇느냐가 다르고, 그래서 타입도 다릅니다.
| 큐 | 타입 | 잇는 경계 | 특징 |
|---|---|---|---|
client_stream |
asyncio.Queue |
코루틴 ↔ 이벤트 루프 | 스레드 안전하지 않음. run_coroutine_threadsafe로만 접근 |
incoming_streaming_queue |
queue.Queue |
스레드 ↔ 스레드 | WorkloadManager의 대기열 |
task_queue / result_queue |
mp.Queue |
프로세스 ↔ 프로세스 | 여기서만 데이터가 복사되어 건너간다 |
asyncio.Queue를 쓰는 이유는 세 가지입니다.
- 대기를 표현하는 언어가 됩니다.
await queue.get()한 줄이 "올 때까지 잠들기"를 그대로 뜻합니다. 큐를 쓰지 않고 같은 일을 하려면Event를 만들어 신호를 주고받거나 폴링 루프를 직접 짜야 하는데, 그러면 코드가 훨씬 길어집니다. - 버퍼 역할을 합니다. 워커는 클라이언트가 느리더라도 기다려주지 않고 곧장 다음 배치로 넘어갑니다. 클라이언트가 아직 읽어가지 못한 토큰은 그동안 큐에 쌓여 있습니다.
- 요청끼리 격리해 줍니다. 요청마다 큐를 따로 만들기 때문에, 배치에 네 개가 함께 실려 있어도 각 시퀀스의 토큰이 서로 섞이지 않습니다.
executor와 worker의 왕복
processing_loop는 큐를 직접 만지지 않습니다. ModelExecutor가 큐 두 개를 소유하고, 건너편에서는 ModelWorker.run()이 무한 루프를 돌며 받아갑니다.
튜플에 실린 True가 저쪽 분기를 결정하고, 돌아온 태그가 이쪽 검증을 통과시킵니다. 프로세스 경계를 오가는 건 이 튜플 두 개가 전부입니다.
여기서 중요한 사실 하나. 워커는 상태가 없습니다. 매번 배치를 통째로 받아 토큰 하나를 돌려주고 잊습니다. 시퀀스가 어디까지 왔는지는 전적으로 WorkloadManager 쪽 기억이고, 그래서 프롬프트를 매 바퀴 다시 실어 보내야 합니다.
합승: 남남인 요청 넷이 행렬 하나에
request_id가 끝까지 따라다녀야 하는 이유가 여기 있습니다. 가운데에서 네 요청은 구분이 사라진 채 하나의 텐서가 되고, 나올 때 순서만으로 다시 갈라집니다.
model()은 "다음 토큰"을 주지 않습니다. 모든 위치에 대해 동시에 "이 자리 다음에 올 토큰의 점수"를 내놓습니다. 학습 때는 그 T개 예측을 전부 쓰지만(teacher forcing), 추론 때는 맨 마지막 하나만 필요해서 [:, -1, :]로 집습니다.
"forward 1회 = 토큰 1개"의 정체가 이겁니다. 모델이 1개만 만들어서가 아니라, 7개분을 계산해놓고 1개만 쓰기 때문입니다.
그리고 결정적으로, generate_forward_batch()의 for i, prompt_data in enumerate(prompts) 루프 안에는 모델 호출이 없습니다. 연산은 이미 끝났고, 그 루프는 [4] 텐서를 인덱스로 풀어 dict 네 개로 포장할 뿐입니다.
추론의 병목은 연산이 아니라 가중치를 HBM에서 읽어오는 대역폭입니다. opt-125m의 가중치를 끌어오는 비용은 배치가 1이든 4든 동일합니다. 그 한 번 읽어온 가중치로 4개를 처리하니 처리량이 4배가 되는 것이고, 이것이 LLM 서빙에서 배치가 핵심인 이유입니다.
분배: 시퀀스가 우편함 주소를 들고 다닌다
워커 스레드는 큐를 따로 보관하지 않습니다. Sequence 객체가 누구냐와 어디로 보내냐를 함께 들고 있어서, request_id로 조회하면 배달 주소가 딸려 나옵니다.
seq.id = "8310f5e1" # 누구냐
seq.prompt = "The quick brown" # 다음 입력
seq.client_stream = <asyncio.Queue> # ★ 어디로
seq.loop = <event loop> # ★ 누가 배달
seq.token_count = 3 # 몇 바퀴째
★ 표시한 둘은 요청이 들어올 때 핸들러가 만들어 시퀀스에 맡겨둔 것입니다.
# llm.py 의 event_generator()
queue = asyncio.Queue() # 여기서 만들어서
seq_id = self.workload_manager.add_streaming_request(prompt, queue, loop)
# workload_manager.py 의 add_streaming_request()
sequence = Sequence(request_id, prompt, client_stream, loop) # 여기에 맡겨둔다
그래서 get_sequence(request_id)가 하는 일은 id로 우편함 주소를 조회하는 것입니다.
두 갈래
두 갈래 모두 ①이 먼저 나가고 ②가 나중에 정리합니다. update_sequence_output은 SSE와 무관한 뒷정리이고, 그 반환값은 아무도 받지 않습니다.
큐에 넣은 값이 소켓에 닿기까지
SSE는 텍스트만 실어 나르는 프로토콜이라 dict를 그대로 보낼 수 없습니다.
종료 신호만은 예외입니다. None은 직렬화하지 않고 그대로 큐에 넣습니다. 받는 쪽이 데이터와 구분해야 하기 때문입니다.
| 큐에 넣는 것 | 코루틴이 하는 일 | |
|---|---|---|
| 계속 | '{"token": ...}' 문자열 |
yield → 소켓에 write |
| 종료 | None |
break → 스트림 닫힘 |
되감기: 매 바퀴 프롬프트가 자란다
use_cache=False라서 매 바퀴 전체 문장을 처음부터 다시 계산합니다. 토큰 하나를 얻으려고 매번 문장 전체를 통과시키므로 뒤로 갈수록 한 바퀴가 느려지고, 21번째 바퀴는 첫 바퀴보다 여섯 배 긴 입력을 처리합니다.
vLLM의 KV 캐시가 없애는 낭비가 정확히 이 부분입니다.
말로 풀면
| 단계 | 무슨 일 | 어디서 |
|---|---|---|
| 01 접수 | /generate_stream이 Sequence를 만든다. 전용 asyncio.Queue와 이벤트 루프 핸들을 붙여두고, 핸들러는 await queue.get()에서 멈춘다 |
이벤트 루프 |
| 02 모집 | 워커 스레드가 0.1초마다 대기열을 훑는다. 발견하면 슬롯에 넣고, 다른 요청이 있으면 최대 4개까지 같이 태운다 | 워커 스레드 |
| 03 발송 | 배치를 ModelExecutor.execute_forward_batch에 넘긴다. task_queue를 소유한 건 executor이고, 여기서 (batch, True)가 복사되어 건너간다. _wait_for_result()로 1초마다 워커 생존을 확인하며 기다린다 |
워커 스레드 |
| 04 연산 | ModelWorker.run()의 무한 루프가 꺼내 batch, is_streaming으로 푼다. True니까 generate_forward_batch로 간다 |
워커 프로세스 |
| 05 포장 | [4]를 풀어 dict 네 개를 만들고 ('stream', results)로 태그를 붙여 넣는다. 워커는 곧장 루프 맨 위로 돌아가 아무것도 기억하지 않는다 |
워커 프로세스 |
| 06 분배 | request_id로 원래 시퀀스를 조회한다. 시퀀스가 자기 큐와 루프를 들고 있으니 어디로 보낼지가 여기서 결정된다 |
워커 스레드 |
| 07 전달 | run_coroutine_threadsafe로 그 큐에 JSON을 넣는다. 이벤트 루프가 01에서 잠든 코루틴을 깨우고, yield된 조각이 소켓으로 나간다 |
이벤트 루프 |
| 08 되감기 | 토큰을 프롬프트에 이어붙이고 token_count를 올린다. SSE는 07에서 이미 끝났으므로 이건 순수한 뒷정리다 |
워커 스레드 |
request_id를 따라가면 길을 잃지 않습니다. 경계를 넘을 때마다 데이터의 형태는 계속 바뀝니다. 처음에는 문자열이었다가 dict가 되고, 워커 안에서는 텐서의 한 행이 되었다가, 다시 dict를 거쳐 JSON 문자열로 나갑니다.