[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번 반복합니다.

request_id가 클라이언트에서 이벤트 루프, 워커 스레드, 워커 프로세스를 거쳐 다시 클라이언트로 돌아오는 시퀀스 다이어그램

큐가 세 종류인 이유

코드에 queue라는 이름이 여러 번 나오는데 전부 다른 물건입니다. 무엇과 무엇을 잇느냐가 다르고, 그래서 타입도 다릅니다.

큐 타입 잇는 경계 특징
client_stream asyncio.Queue 코루틴 ↔ 이벤트 루프 스레드 안전하지 않음. run_coroutine_threadsafe로만 접근
incoming_streaming_queue queue.Queue 스레드 ↔ 스레드 WorkloadManager의 대기열
task_queue / result_queue mp.Queue 프로세스 ↔ 프로세스 여기서만 데이터가 복사되어 건너간다

클라이언트, 이벤트 루프, 워커 스레드, 워커 프로세스 네 구역과 그 사이를 잇는 세 종류의 큐를 보여주는 계층도

asyncio.Queue를 쓰는 이유는 세 가지입니다.

  1. 대기를 표현하는 언어가 됩니다. await queue.get() 한 줄이 "올 때까지 잠들기"를 그대로 뜻합니다. 큐를 쓰지 않고 같은 일을 하려면 Event를 만들어 신호를 주고받거나 폴링 루프를 직접 짜야 하는데, 그러면 코드가 훨씬 길어집니다.
  2. 버퍼 역할을 합니다. 워커는 클라이언트가 느리더라도 기다려주지 않고 곧장 다음 배치로 넘어갑니다. 클라이언트가 아직 읽어가지 못한 토큰은 그동안 큐에 쌓여 있습니다.
  3. 요청끼리 격리해 줍니다. 요청마다 큐를 따로 만들기 때문에, 배치에 네 개가 함께 실려 있어도 각 시퀀스의 토큰이 서로 섞이지 않습니다.

executor와 worker의 왕복

processing_loop는 큐를 직접 만지지 않습니다. ModelExecutor가 큐 두 개를 소유하고, 건너편에서는 ModelWorker.run()이 무한 루프를 돌며 받아갑니다.

ModelExecutor와 ModelWorker.run 사이를 task_queue와 result_queue로 두 개의 튜플이 오가는 왕복 구조

튜플에 실린 True가 저쪽 분기를 결정하고, 돌아온 태그가 이쪽 검증을 통과시킵니다. 프로세스 경계를 오가는 건 이 튜플 두 개가 전부입니다.

여기서 중요한 사실 하나. 워커는 상태가 없습니다. 매번 배치를 통째로 받아 토큰 하나를 돌려주고 잊습니다. 시퀀스가 어디까지 왔는지는 전적으로 WorkloadManager 쪽 기억이고, 그래서 프롬프트를 매 바퀴 다시 실어 보내야 합니다.

합승: 남남인 요청 넷이 행렬 하나에

request_id가 끝까지 따라다녀야 하는 이유가 여기 있습니다. 가운데에서 네 요청은 구분이 사라진 채 하나의 텐서가 되고, 나올 때 순서만으로 다시 갈라집니다.

서로 다른 네 개의 요청이 하나의 배치 텐서로 묶여 한 번의 forward를 거친 뒤 다시 네 개의 토큰으로 갈라지는 그림

model()은 "다음 토큰"을 주지 않습니다. 모든 위치에 대해 동시에 "이 자리 다음에 올 토큰의 점수"를 내놓습니다. 학습 때는 그 T개 예측을 전부 쓰지만(teacher forcing), 추론 때는 맨 마지막 하나만 필요해서 [:, -1, :]로 집습니다.

"forward 1회 = 토큰 1개"의 정체가 이겁니다. 모델이 1개만 만들어서가 아니라, 7개분을 계산해놓고 1개만 쓰기 때문입니다.

그리고 결정적으로, generate_forward_batch()의 for i, prompt_data in enumerate(prompts) 루프 안에는 모델 호출이 없습니다. 연산은 이미 끝났고, 그 루프는 [4] 텐서를 인덱스로 풀어 dict 네 개로 포장할 뿐입니다.

무거운 연산은 나눠 쓰고, 파이썬 루프만 4번 돈다

추론의 병목은 연산이 아니라 가중치를 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로 우편함 주소를 조회하는 것입니다.

두 갈래

워커 결과 하나가 request_id로 시퀀스를 조회한 뒤, 종료 여부에 따라 두 갈래로 나뉘어 처리되는 흐름

두 갈래 모두 ①이 먼저 나가고 ②가 나중에 정리합니다. update_sequence_output은 SSE와 무관한 뒷정리이고, 그 반환값은 아무도 받지 않습니다.

큐에 넣은 값이 소켓에 닿기까지

SSE는 텍스트만 실어 나르는 프로토콜이라 dict를 그대로 보낼 수 없습니다.

파이썬 딕셔너리가 JSON 문자열로 직렬화되고 SSE 프레임으로 감싸여 소켓에 나가는 3단계

종료 신호만은 예외입니다. 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 문자열로 나갑니다.