LabHub
시작하기
배우기 러닝패스 코스

실시간 통신 — WebSocket·gRPC 스트리밍·WebRTC

오래 붙어 있는 WebSocket 허브를 운영한다

LabHub 에서 이어서 보기

목표

websockets 로 발행·구독 허브와 재접속 클라이언트를 만들어, 느린 구독자·유휴 타임아웃·말없이 사라진 상대·배포로 인한 일괄 끊김을 차례로 견디게 합니다.

왜 중요한가

WebSocket 서비스의 장애는 대부분 프로토콜이 아니라 시간에서 옵니다. 한 명이 느려지면 허브 메모리가 차고, 조용한 연결은 로드밸런서가 먼저 끊고, 휴대폰이 터널에 들어가면 상대는 close 없이 사라지고, 배포 한 번에 수만 개가 동시에 다시 붙습니다. 이 실습의 다섯 규칙 — 구독자별 상한, 1013 으로 끊기, 유휴 한도보다 짧은 ping, 닫힘을 함께 기다리기, 지터가 있는 재접속과 seq 로 이어 받기 — 는 채팅이든 시세든 음성 AI 의 부분 결과든 똑같이 필요합니다.

단계

  1. 재접속 간격에 지터를 섞는다 — /root/rt/wsops/client.py 에 backoff(attempt, base=0.2, cap=5.0, rng=random) 을 만드세요. attempt 번째 재시도 전에 기다릴 초를 돌려줍니다. 값은 rng.uniform(0, min(cap, base * 2 ** attempt)) 한 번으로 정합니다(full jitter). rng 에는 random 모듈이나 random.Random 객체가 옵니다.
  2. 닫기 코드를 보고 다시 붙을지 정한다 — /root/rt/wsops/client.py 에 should_reconnect(code) 를 추가하세요. 1001(가는 중)·1006(close 프레임 없이 끊김)·1011(서버 내부 오류)·1012(서비스 재시작)·1013(나중에 다시)이면 True, 그 밖의 코드는 모두 False 를 돌려줍니다.
  3. 허브가 한 사람의 말을 모두에게 나른다 — /root/rt/wsops/hub.py 에 run(host, port, ring=1000, queue=100, ping_interval=1.0, ping_timeout=2.0) 코루틴을 만드세요. websockets.asyncio.server.serve 로 경로 셋을 받습니다. /pub 로 들어온 텍스트 메시지마다 1부터 늘어나는 seq 를 붙여 {"seq": n, "data": 메시지} JSON 문자열로 만들고, /sub 로 붙은 모든 구독자에게 보냅니다. /stats 는 {"subscribers": 구독자 수, "dropped": 끊어 낸 구독자 수, "seq": 마지막 seq} JSON 하나를 보내고 끝냅니다. 구독자마다 asyncio.Queue 와 그 큐를 비우며 보내는 루프를 따로 둡니다.
  4. 느린 구독자 한 명만 끊는다 — 구독자 큐의 크기를 queue 로 제한하세요. 큐가 가득 차 넣을 수 없으면 그 구독자를 목록에서 빼고 dropped 를 1 올린 뒤 닫기 코드 1013 과 이유 "slow consumer" 로 닫습니다. 채점기는 읽지 않는 구독자 하나와 잘 읽는 구독자 하나를 붙인 채 8000바이트 메시지 8000개를 발행하고, 잘 읽는 쪽이 전부 받는지, 허브의 메모리가 얼마나 늘었는지, 읽지 않던 쪽이 1013 을 받는지 봅니다.
  5. 조용한 연결이 로드밸런서에 잘리지 않게 한다 — serve 에 ping_interval 과 ping_timeout 을 run 의 인자 그대로 넘기세요. 채점기는 2.5초 동안 조용한 연결을 끊는 중계기를 허브 앞에 두고, ping 을 스스로 보내지 않는 구독자로 6초 동안 기다린 뒤 메시지 하나를 발행합니다. 그 메시지가 구독자에게 닿아야 합니다.
  6. 말없이 사라진 구독자를 치운다 — 구독자 루프가 큐만 기다리지 말고 연결이 닫히는 것도 함께 기다리게 고치세요. ws.wait_closed() 와 q.get() 을 asyncio.wait(..., return_when=FIRST_COMPLETED) 로 함께 기다리다가 연결이 먼저 닫히면 구독자 목록에서 빼고 루프를 끝냅니다. serve 에는 close_timeout=1.0 도 넘기세요. 채점기는 핸드셰이크만 하고 ping 에 답하지 않는 구독자를 붙인 뒤, 6초 안에 /stats 의 subscribers 가 0 으로 돌아오는지 봅니다.
  7. 끊긴 자리부터 이어 받는다 — 허브에 최근 ring 개의 메시지를 담는 링 버퍼를 두고, /sub?last=N 으로 붙으면 seq 가 N 보다 큰 메시지를 먼저 다시 보낸 뒤 실시간으로 이어 보내게 하세요. N+1 이 버퍼에서 이미 밀려났으면 {"reset": true, "seq": 마지막 seq} 를 먼저 보냅니다. 그리고 /root/rt/wsops/client.py 에 consume(url, out_path, stop_after, base=0.1, cap=1.0) 코루틴을 추가하세요. url 에 붙어 받은 메시지의 seq 를 한 줄에 하나씩 out_path 에 덧붙이고, 연결이 끊기면 should_reconnect 와 backoff 로 기다렸다가 마지막으로 적은 seq 를 last 로 넘겨 다시 붙습니다. seq 가 stop_after 에 닿으면 끝냅니다. 채점기는 발행 도중 연결을 두 번 끊고, 파일에 1 부터 stop_after 까지가 빠짐없이 한 번씩 있는지 봅니다.

참고

재접속 간격에 지터를 섞는다

/root/rt/wsops/client.py 에 backoff(attempt, base=0.2, cap=5.0, rng=random) 을 만드세요. attempt 번째 재시도 전에 기다릴 초를 돌려줍니다. 값은 rng.uniform(0, min(cap, base * 2 ** attempt)) 한 번으로 정합니다(full jitter). rng 에는 random 모듈이나 random.Random 객체가 옵니다.

지터가 없으면 한꺼번에 끊긴 클라이언트 수만 개가 정확히 같은 시각에 다시 몰려와 막 살아난 서버를 다시 쓰러뜨립니다. 상한(cap)이 없으면 열 번째 재시도는 몇 분 뒤가 됩니다. 채점기는 씨앗을 정한 random.Random 을 넘겨 같은 값을 계산해 봅니다.

닫기 코드를 보고 다시 붙을지 정한다

/root/rt/wsops/client.py 에 should_reconnect(code) 를 추가하세요. 1001(가는 중)·1006(close 프레임 없이 끊김)·1011(서버 내부 오류)·1012(서비스 재시작)·1013(나중에 다시)이면 True, 그 밖의 코드는 모두 False 를 돌려줍니다.

1008(정책 위반)이나 1002(프로토콜 오류)로 끊긴 연결에 다시 붙으면 똑같은 이유로 또 끊깁니다. 재시도가 의미 있는 것은 상대 쪽 사정이 바뀔 수 있는 경우뿐입니다. 1000 은 누군가 일부러 닫았다는 뜻입니다.

허브가 한 사람의 말을 모두에게 나른다

/root/rt/wsops/hub.py 에 run(host, port, ring=1000, queue=100, ping_interval=1.0, ping_timeout=2.0) 코루틴을 만드세요. websockets.asyncio.server.serve 로 경로 셋을 받습니다. /pub 로 들어온 텍스트 메시지마다 1부터 늘어나는 seq 를 붙여 {"seq": n, "data": 메시지} JSON 문자열로 만들고, /sub 로 붙은 모든 구독자에게 보냅니다. /stats 는 {"subscribers": 구독자 수, "dropped": 끊어 낸 구독자 수, "seq": 마지막 seq} JSON 하나를 보내고 끝냅니다. 구독자마다 asyncio.Queue 와 그 큐를 비우며 보내는 루프를 따로 둡니다.

발행하는 쪽이 구독자마다 await send 를 차례로 부르면, 한 명이 느릴 때 그 뒤의 모두가 기다립니다. 구독자별 큐는 그 기다림을 사람마다 떼어 놓는 장치입니다. 발행은 큐에 넣기만 하고 기다리지 않습니다.

느린 구독자 한 명만 끊는다

구독자 큐의 크기를 queue 로 제한하세요. 큐가 가득 차 넣을 수 없으면 그 구독자를 목록에서 빼고 dropped 를 1 올린 뒤 닫기 코드 1013 과 이유 "slow consumer" 로 닫습니다. 채점기는 읽지 않는 구독자 하나와 잘 읽는 구독자 하나를 붙인 채 8000바이트 메시지 8000개를 발행하고, 잘 읽는 쪽이 전부 받는지, 허브의 메모리가 얼마나 늘었는지, 읽지 않던 쪽이 1013 을 받는지 봅니다.

큐에 상한이 없으면 느린 구독자 한 명 몫의 메시지가 허브 메모리에 끝없이 쌓입니다. 메시지를 버리는 것과 연결을 끊는 것 중 무엇이 나은지는 데이터의 성격이 정합니다. 순서와 누락이 중요한 스트림이라면 구멍 난 채로 계속 보내느니 끊고 다시 붙게 하는 편이 정직합니다.

조용한 연결이 로드밸런서에 잘리지 않게 한다

serve 에 ping_interval 과 ping_timeout 을 run 의 인자 그대로 넘기세요. 채점기는 2.5초 동안 조용한 연결을 끊는 중계기를 허브 앞에 두고, ping 을 스스로 보내지 않는 구독자로 6초 동안 기다린 뒤 메시지 하나를 발행합니다. 그 메시지가 구독자에게 닿아야 합니다.

websockets 의 기본 ping 간격은 20초라, 유휴 한도가 그보다 짧은 장비 뒤에서는 조용한 연결이 먼저 끊깁니다. 흔한 AWS ALB 의 기본 유휴 한도는 60초이고 사내 프록시는 더 짧은 경우가 많습니다. 간격은 가장 짧은 유휴 한도보다 짧아야 합니다.

말없이 사라진 구독자를 치운다

구독자 루프가 큐만 기다리지 말고 연결이 닫히는 것도 함께 기다리게 고치세요. ws.wait_closed() 와 q.get() 을 asyncio.wait(..., return_when=FIRST_COMPLETED) 로 함께 기다리다가 연결이 먼저 닫히면 구독자 목록에서 빼고 루프를 끝냅니다. serve 에는 close_timeout=1.0 도 넘기세요. 채점기는 핸드셰이크만 하고 ping 에 답하지 않는 구독자를 붙인 뒤, 6초 안에 /stats 의 subscribers 가 0 으로 돌아오는지 봅니다.

ping 시간 초과로 라이브러리가 연결을 닫아도, 큐를 기다리는 코루틴은 새 메시지가 오기 전까지 깨지 않습니다. 그 사이 구독자 목록에는 죽은 연결이 남아 숫자를 부풀리고 메시지를 받아 쌓습니다. 그리고 답 없는 상대에게 close 를 보낸 뒤 답을 기다리는 시간의 기본값은 10초입니다.

끊긴 자리부터 이어 받는다

허브에 최근 ring 개의 메시지를 담는 링 버퍼를 두고, /sub?last=N 으로 붙으면 seq 가 N 보다 큰 메시지를 먼저 다시 보낸 뒤 실시간으로 이어 보내게 하세요. N+1 이 버퍼에서 이미 밀려났으면 {"reset": true, "seq": 마지막 seq} 를 먼저 보냅니다. 그리고 /root/rt/wsops/client.py 에 consume(url, out_path, stop_after, base=0.1, cap=1.0) 코루틴을 추가하세요. url 에 붙어 받은 메시지의 seq 를 한 줄에 하나씩 out_path 에 덧붙이고, 연결이 끊기면 should_reconnect 와 backoff 로 기다렸다가 마지막으로 적은 seq 를 last 로 넘겨 다시 붙습니다. seq 가 stop_after 에 닿으면 끝냅니다. 채점기는 발행 도중 연결을 두 번 끊고, 파일에 1 부터 stop_after 까지가 빠짐없이 한 번씩 있는지 봅니다.

"받은 것" 이 아니라 "처리해서 적은 것" 을 기준으로 다시 붙어야 합니다. 받기만 하고 적기 전에 끊기면 그 메시지는 영영 사라집니다. 같은 seq 가 두 번 오면 한 번만 적으세요.