본문 바로가기
C.W.K.
Stream
Lesson 07 of 07 · published

Redis Pub/Sub 연결 다리

~14 min · management, scaling, redis, pubsub

Level 0Poller
0 XP0/60 lessons0/10 achievements
0/120 XP to next level120 XP to go0% complete

서버 한 대의 한계

FastAPI 프로세스 하나가 WebSocket 연결 수만 개를 다룰 수는 있지만 백만 개까지는 어려워. 부하 분산기 뒤에 서버를 여러 대 두고 수평 확장하면 구조적인 문제가 생겨. 같은 방의 연결이 서로 다른 서버에 있을 수 있어서 메시지 버스를 공유하지 않으면 서버 A 의 전체 전송이 서버 B 의 사용자에게 닿지 않아.

표준 해법은 Redis Pub/Sub 이야

각 서버는 나가는 메시지를 Redis 채널에 발행하고 같은 채널을 구독해. Redis 가 모든 구독자에게 메시지를 펼쳐 주면 각 서버는 자기 로컬 연결에만 다시 전송하지. 중앙에 연결 상태를 두지 않고도 이 패턴은 서버 수천 대와 연결 수백만 개까지 확장할 수 있어.

cwkPippa 의 비슷한 패턴

cwkPippa 는 집의 Mac Studio 한 대에서 의도적으로 단일 서버로 돌아가. 하지만 사용자의 질문 하나를 여러 브레인 프로세스에 펼쳤다가 응답을 모으는 피파의 Council 패턴은 같은 모양이야. Redis 대신 하위 프로세스 파이프를 쓴다는 운반 방식만 다를 뿐, 펼치고 모으는 구조는 분산 실시간 시스템 전반에 나타나.

Code

Redis 를 사용하는 연결 관리자·python
import json
import redis.asyncio as aioredis
from fastapi import FastAPI, WebSocket
from typing import Dict, Set
import asyncio

class RedisManager:
    def __init__(self, redis_url: str):
        self.redis = aioredis.from_url(redis_url)
        self.pubsub = self.redis.pubsub()
        self.local: Dict[str, Set[WebSocket]] = {}
        self._task: asyncio.Task | None = None

    async def start(self):
        await self.pubsub.psubscribe('ws:room:*')
        self._task = asyncio.create_task(self._listen())

    async def stop(self):
        if self._task:
            self._task.cancel()
        await self.pubsub.aclose()
        await self.redis.aclose()

    async def join(self, ws: WebSocket, room: str):
        await ws.accept()
        self.local.setdefault(room, set()).add(ws)

    def leave(self, ws: WebSocket, room: str):
        members = self.local.get(room)
        if members:
            members.discard(ws)
            if not members:
                self.local.pop(room, None)

    async def publish(self, room: str, message: dict):
        '''All servers (including this one) will fan it out locally.'''
        await self.redis.publish(f'ws:room:{room}', json.dumps(message))

    async def _listen(self):
        async for msg in self.pubsub.listen():
            if msg.get('type') != 'pmessage':
                continue
            channel = msg['channel'].decode()
            room = channel.split(':', 2)[2]
            data = json.loads(msg['data'])
            for ws in list(self.local.get(room, ())):
                try:
                    await ws.send_json(data)
                except Exception:
                    self.leave(ws, room)
구성도·text
  Server A              Redis              Server B
  (A, B in #general)   pub/sub            (C, D in #general)
        |                |                       |
  user A sends "hi"      |                       |
        |--- publish --->|                       |
        |                |--- pmessage --------->|
        |                |                       v
        |                |                C and D get "hi"
        v                |
  A and B get "hi"
  (via local fan-out)

External links

Exercise

FastAPI 서버 두 대를 8001 과 8002 포트에 띄우고 같은 Redis 에 연결해. 각각 클라이언트를 연 뒤 8001 쪽에서 메시지를 보내 8002 쪽이 Redis 연결 다리를 통해 받는지 확인해. Redis 를 멈추면 서버 간 전송도 멈추는지, 다시 시작하면 재개되는지도 확인해.

Progress

Progress is local-only — sign in to sync across devices.
이 페이지에서 버그를 발견하셨거나 피드백이 있으세요?문제 신고

댓글 0

🔔 답글 알림 (로그인 필요)
로그인댓글을 남기려면 로그인해 주세요.

아직 댓글이 없어요. 첫 댓글을 남겨보세요.