[CloudNeta] Hands-On LLM Serving 2주차 - 토끼굴: 파이프와 Connection 계층
이 문서는 2주차 part 1 - 워커 프로세스는 왜 '서 있는가'에서 갈라져 나온 토끼굴입니다. 요약하면 파이프는 부모-자식 간 통신하는 통로이고, Queue 내부에서 통신 채널은 이렇게 만들어집니다.
self._reader, self._writer = connection.Pipe(duplex=False)
파이프란
파이프는 두 프로세스가 통신하는 전달자로, UNIX 계열에서 쓰이는 IPC(Inter-Process Communication) 방식 중 하나입니다[1]. duplex=False로 만든 파이프는 단방향입니다. 한쪽 끝(_writer)은 쓰기만, 반대쪽 끝(_reader)은 읽기만 합니다. 그래서 부모와 워커가 왕복 통신을 하려면 파이프가 두 개 필요하고, 이 코드베이스의 task_queue(부모→워커)와 result_queue(워커→부모)가 정확히 그 두 개입니다.
파이썬 문서가 약속하는 계약
multiprocessing 공식 문서의 Connection.recv_bytes() 설명입니다[2].
접속의 반대편 끝에서 송신된 바이트 데이터의 완전한 메시지를 돌려줍니다. 뭔가 수신할 때까지 블록합니다. 수신할 내용이 없고 반대편 끝이 닫혔으면 EOFError를 발생시킵니다.
"완전한 메시지"와 "수신할 때까지 블록"이라는 두 약속이 핵심입니다. 아래에서 이 약속이 코드로 어떻게 구현되는지 확인합니다.
인터페이스와 플랫폼별 구현
multiprocessing/connection.py는 인터페이스와 구현이 분리되어 있습니다[3].
_ConnectionBase.recv_bytes(): 공용 인터페이스. 닫힘/읽기 가능 여부를 검사한 뒤self._recv_bytes()에 위임합니다.PipeConnection._recv_bytes()(윈도우 네이티브 전용):if _winapi:아래에 정의되며, named pipe에 대해ReadFile(overlapped=True)후WaitForMultipleObjects(..., INFINITE)로 대기합니다. 리눅스의 커널 수면과 역할이 같은 윈도우식 블로킹입니다.Connection._recv_bytes()(유닉스):_read = os.read로 묶여 있고,_recv()루프에서chunk = read(handle, remaining)을 호출합니다. 이read가 운영체제의read(2)이며, 파이썬의 종점입니다.
connection.py 소스 발췌 (주석 원문 유지)
# multiprocessing.connection.py
######### 인터페이스는 이거고
class _ConnectionBase:
...
def recv_bytes(self, maxlength=None):
"""
Receive bytes data as a bytes object.
"""
self._check_closed()
self._check_readable()
if maxlength is not None and maxlength < 0:
raise ValueError("negative maxlength")
buf = self._recv_bytes(maxlength)
if buf is None:
self._bad_message_length()
return buf.getvalue()
######### 윈도우즈 구현체는 이거고
if _winapi:
class PipeConnection(_ConnectionBase):
...
def _recv_bytes(self, maxsize=None):
if self._got_empty_message:
self._got_empty_message = False
return io.BytesIO()
else:
bsize = 128 if maxsize is None else min(maxsize, 128)
try:
ov, err = _winapi.ReadFile(self._handle, bsize,
overlapped=True)
try:
if err == _winapi.ERROR_IO_PENDING:
waitres = _winapi.WaitForMultipleObjects(
[ov.event], False, INFINITE)
assert waitres == WAIT_OBJECT_0
except:
ov.cancel()
raise
finally:
nread, err = ov.GetOverlappedResult(True)
if err == 0:
f = io.BytesIO()
f.write(ov.getbuffer())
return f
elif err == _winapi.ERROR_MORE_DATA:
return self._get_more_data(ov, maxsize)
except OSError as e:
if e.winerror == _winapi.ERROR_BROKEN_PIPE:
raise EOFError
else:
raise
raise RuntimeError("shouldn't get here; expected KeyboardInterrupt")
########## xNix 구현체는 이거다.
########## 정확히는 유닉스에서 파이프/소켓 공용 구현체고,
########## 윈도우에서는 소켓 기반 연결에만 쓰인다 (윈도우의 Pipe()는 위의 PipeConnection을 씀)
class Connection(_ConnectionBase):
...
if _winapi:
def _close(self, _close=_multiprocessing.closesocket):
_close(self._handle)
_write = _multiprocessing.send
_read = _multiprocessing.recv
else:
def _close(self, _close=os.close):
_close(self._handle)
_write = os.write
# 유닉스면 os.read다. WSL도 리눅스라서 이 길을 탄다
_read = os.read
def _recv(self, size, read=_read):
buf = io.BytesIO()
handle = self._handle
remaining = size
while remaining > 0:
# 이 read가 운영체제의 read다. read(2) 시스템콜.
# 빈 파이프면 여기서 커널이 프로세스를 재운다. "서 있는" 지점의 실체.
chunk = read(handle, remaining)
n = len(chunk)
if n == 0:
if remaining == size:
raise EOFError
else:
raise OSError("got end of file during message")
buf.write(chunk)
remaining -= n
return buf
def _recv_bytes(self, maxsize=None):
# 이 함수는 운영체제 함수가 아니라 메시지 프레이밍 계층이다.
# os.read는 "최대 N바이트"만 보장하고 메시지 경계를 모르니,
# 4바이트 길이 헤더(-1이면 8바이트 확장)를 먼저 읽고 본문을 그만큼 채워서
# 문서가 약속한 "완전한 메시지 하나"를 복원한다
buf = self._recv(4)
size, = struct.unpack("!i", buf.getvalue())
if size == -1:
buf = self._recv(8)
size, = struct.unpack("!Q", buf.getvalue())
if maxsize is not None and size > maxsize:
return None
return self._recv(size)
처음에 틀렸던 것: WSL은 윈도우가 아니다
처음에는 "WSL이니까 윈도우 경로(PipeConnection, AF_PIPE)를 타겠지"라고 추측했는데, 반대였습니다. WSL2는 진짜 리눅스 커널 위에서 도는 리눅스 환경이고, 그 안의 파이썬은 리눅스 빌드라서 _winapi 모듈 자체가 없습니다. AF_PIPE는 윈도우 네이티브 파이썬 전용 주소 패밀리입니다. 증거는 세 가지입니다.
ps에서 관측된wchan=pipe_read는 리눅스 커널fs/pipe.c의 함수명입니다. 윈도우 named pipe였다면 나올 수 없습니다.- 이 프로젝트는
multiprocessing.get_context("fork")를 쓰는데, fork는 윈도우에 존재하지 않는 시작 방식입니다. 서버가 뜬다는 사실 자체가 유닉스 경로의 증명입니다. - WSL 셸에서
python -c "import _winapi"를 실행하면ModuleNotFoundError가 납니다.
llm-study/ch03/single_model_study on main [?] via 🐍 v3.12.12
➜ python -c "import _winapi"
Traceback (most recent call last):
File "<string>", line 1, in <module>
ModuleNotFoundError: No module named '_winapi'
_recv_bytes()는 OS 함수가 아니라 프레이밍 계층이다
os.read는 "바이트 스트림에서 최대 N바이트 읽기"만 보장할 뿐 메시지 경계를 모릅니다. 문서가 약속한 "완전한 메시지"를 복원하는 것이 _recv_bytes()의 역할입니다.
def _recv_bytes(self, maxsize=None):
buf = self._recv(4) # 4바이트 길이 헤더 (big-endian)
size, = struct.unpack("!i", buf.getvalue())
if size == -1: # 4바이트로 표현 못 하는 큰 메시지
buf = self._recv(8) # 8바이트 확장 길이를 다시 읽음
size, = struct.unpack("!Q", buf.getvalue())
...
return self._recv(size) # 본문을 정확히 size만큼 채워서 반환
_recv()는 요청한 크기가 다 찰 때까지 os.read를 반복하는 루프입니다. 그래서 큐가 비어 있을 때 잠드는 지점은 정확히 "첫 4바이트 헤더를 읽으려는 os.read"가 됩니다.
파이프 입문 글(sikpang). 파이프 개념을 쉽게 정리한 한국어 자료. ↩︎
multiprocessing 공식 문서(한국어).
Connection.send_bytes()/recv_bytes()의 계약. ↩︎CPython Lib/multiprocessing/connection.py.
_ConnectionBase와 플랫폼별_recv_bytes()구현. ↩︎