메신저 만들기
Hermaeus Mora · · Backend
스키마
CREATE TABLE room_seq (
room_id BIGINT PRIMARY KEY,
seq BIGINT NOT NULL DEFAULT 0
);
CREATE TABLE messages (
room_id BIGINT NOT NULL,
seq BIGINT NOT NULL,
client_msg_id TEXT NOT NULL,
sender_id BIGINT NOT NULL,
body TEXT NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
PRIMARY KEY (room_id, seq),
UNIQUE (room_id, client_msg_id) -- 재전송 멱등성의 실제 강제 지점
);seq 카운터가 메시지와 같은 내구성 도메인에 있는 게 핵심. 둘이 같이 죽고 같이 살아나니 어긋날 방법이 없고, 대조·복구 로직이 아예 필요 없다. 카운터를 Redis 같은 휘발성 저장소에 두면 유실 시 seq가 1부터 재시작하는데, INCR은 없는 키에 조용히 1을 반환하므로 에러 없이 번호가 재사용된다. 그 순간 "내 마지막 seq 이후 주세요" 복구 경로가 통째로 깨진다.
라이프사이클
room_locks: defaultdict[int, asyncio.Lock] = defaultdict(asyncio.Lock)
async def join_room(ws, room_id: int):
async with room_locks[room_id]:
first_local = not local_rooms[room_id]
local_rooms[room_id].add(ws)
if first_local:
# 이 프로세스에 이 방 사람이 처음 생겼을 때만 Redis 구독
await pubsub.subscribe(f"room:{room_id}")
async def leave_room(ws, room_id: int):
async with room_locks[room_id]:
local_rooms[room_id].discard(ws)
if not local_rooms[room_id]:
await pubsub.unsubscribe(f"room:{room_id}")
if not local_rooms[room_id]: # await 사이에 재입장했는지 확인
del local_rooms[room_id]
else:
await pubsub.subscribe(f"room:{room_id}")락과 재확인이 필요한 이유. del을 await unsubscribe보다 먼저 하면, 그 await 도중 같은 방에 재입장이 들어왔을 때 local_rooms엔 소켓이 있는데 구독은 끊긴 좀비 방이 생긴다. 프로세스는 겉보기 멀쩡하고 그 방 메시지만 안 오는, 제일 찾기 어려운 형태의 장애다.
클라이언트 송신
clientMsgId 생성 → 화면에 status: "sending"으로 즉시 렌더(optimistic) → WS로 전송.
아직 seq가 없으므로 확정 메시지들과 같은 리스트에 섞지 않는다. 로컬 pending은 별도 tail 영역에 고정해서 쌓고, 서버 seq를 받은 뒤에 본 리스트로 옮긴다. 그냥 "맨 아래"에 두면 내가 보내는 도중 남의 확정 메시지가 도착할 때 위치가 튄다.
ws-a, ws-b 파드
채번·저장·팬아웃을 이 순서로 한다. 순서가 곧 보장의 순서다.
async def handle_send(room_id: int, sender_id: int, client_msg_id: str, body: str):
# 0) 재전송 확인 — 채번 전에 본다. 여기서 걸러야 seq가 낭비되지 않는다
seq = await db.fetchval(
"SELECT seq FROM messages WHERE room_id = $1 AND client_msg_id = $2",
room_id, client_msg_id,
)
if seq is None:
# 1) 채번 + 저장을 한 트랜잭션으로. 여기까지 커밋되면 메시지는 살아남는다
async with db.transaction():
seq = await db.fetchval("""
INSERT INTO room_seq (room_id, seq) VALUES ($1, 1)
ON CONFLICT (room_id) DO UPDATE SET seq = room_seq.seq + 1
RETURNING seq
""", room_id)
await db.execute("""
INSERT INTO messages (room_id, seq, client_msg_id, sender_id, body)
VALUES ($1, $2, $3, $4, $5)
ON CONFLICT (room_id, client_msg_id) DO NOTHING
""", room_id, seq, client_msg_id, sender_id, body)
payload = {"roomId": room_id, "seq": seq, "clientMsgId": client_msg_id,
"senderId": sender_id, "body": body}
# 2) 부가 처리 팬아웃. 실패해도 메시지는 DB에 있으므로 복구 경로가 메꾼다
await producer.send_and_wait("chat.messages", key=room_id, value=payload)
# 3) 실시간 팬아웃
await redis.publish(f"room:{room_id}", payload)
return payload채번을 UPDATE ... WHERE가 아니라 UPSERT로 하는 이유. 신규 방은 room_seq 행이 없어서 UPDATE가 0행을 갱신하고 RETURNING이 빈 결과를 준다.
트랜잭션 경계. 커밋 → Kafka send → publish 순서다. Kafka ack를 트랜잭션 안에서 기다리면 room_seq 행 잠금을 Kafka 왕복(515ms) 동안 붙들어서 방 처리량이 초당 60200으로 떨어진다. 커밋을 먼저 끝내면 잠금 보유 시간이 DB 쓰기 시간(~0.5ms)으로 줄고, Kafka send가 실패해도 메시지는 이미 내구성이 확보돼 있어 갭이 생기지 않는다.
지연 비용. DB 커밋 + acks="all" 왕복이 모두 사용자 대기 시간에 들어간다. 합쳐서 대략 6~20ms. 정확성과 맞바꾼 값이고, 사람이 채팅을 치는 속도 기준으로는 문제되지 않는다.
처리량 한계. room_seq 행 잠금은 커밋까지 유지되므로 방 하나의 처리량은 대략 1 / 트랜잭션 시간 = 초당 500~2000건. 사람이 쓰는 방은 여유가 1000배다. 전역 쓰기량이 초당 1만 건을 넘어가면 room_id로 샤딩한다. 방 하나는 샤드 하나에 사니 방 단위 순서는 구조적으로 보존된다.
프로듀서 설정.
acks="all" # ISR 전체 반영 확인
enable_idempotence=True # 재시도 시 파티션 내 중복/재정렬 방지
max_in_flight_requests_per_connection=5 # 멱등 활성 시 5까지 순서 안전key=room_id가 보장하는 건 "같은 방은 한 파티션"까지다. 파티션 안의 순서는 브로커 도착 순서라서 서로 다른 파드가 같은 방에 동시 전송하면 seq 순서와 어긋날 수 있다. 순서의 진실은 항상 seq이고 컨슈머가 그 기준으로 정렬한다. 파티션 수 N은 방 개수와 무관하게 처리량 기준으로 고정.
Redis pub/sub 팬아웃
room:5를 구독 중인 ws-a, ws-b 파드의 pubsub_reader 태스크가 메시지를 받아 local_rooms[5]의 소켓들에 asyncio.gather로 병렬 push → Carol, Dave, Emma 즉시 수신.
Bob에게도 반사돼 돌아오고, 클라이언트는 clientMsgId로 기존 껍데기를 찾아 서버 seq를 덮어쓰고 status: "sent"로 바꾼다.
이 태스크가 죽으면 프로세스는 겉보기 멀쩡한데 메시지만 안 오는 상태가 되므로, 죽으면 반드시 되살아나야 한다.
def spawn_pubsub_reader():
task = asyncio.create_task(pubsub_reader())
task.add_done_callback(on_reader_dead)
return task
def on_reader_dead(task: asyncio.Task):
if shutting_down.is_set() or task.cancelled():
return
log.error("pubsub_reader died", exc_info=task.exception())
asyncio.get_running_loop().call_later(1.0, resubscribe_and_respawn)
async def resubscribe_and_respawn():
await pubsub.reset()
for room_id in list(local_rooms.keys()): # 이 파드가 들고 있던 방 전부 재구독
await pubsub.subscribe(f"room:{room_id}")
spawn_pubsub_reader()ws 파드 SIGTERM
terminationGracePeriodSeconds: 300
readinessProbe:
httpGet: { path: /readyz, port: 8080 } # `/readyz`는 `shutting_down.is_set()`이면 503을 반환
periodSeconds: 2
failureThreshold: 1 # 즉시 NotReady. 기본값 3이면 6초를 더 기다린다
strategy:
rollingUpdate:
maxSurge: 2
maxUnavailable: 0 # 새 파드가 Ready된 다음에 구 파드를 빼야 재접속을 받아줌import signal, asyncio, random
shutting_down = asyncio.Event()
async def graceful_shutdown():
if shutting_down.is_set(): # SIGTERM이 두 번 와도 태스크는 하나만
return
shutting_down.set() # readinessProbe가 이제 503 반환
# 1) LB 타겟 등록 해제 전파 대기
# Endpoints 갱신 + kube-proxy/LB 반영까지. ALB면 deregistration_delay도 별도로 있다
# 이 동안에도 기존 커넥션은 정상 동작한다
await asyncio.sleep(25)
# 2) 커넥션을 천천히 빼기 (초당 약 1000소켓)
all_ws = [ws for socks in local_rooms.values() for ws in socks]
random.shuffle(all_ws)
for i, ws in enumerate(all_ws):
await ws.send_json({"type": "reconnect", "after_ms": random.randint(0, 60_000)})
await ws.close(code=1001)
if i % 500 == 0:
await asyncio.sleep(0.5)
# 3) 정리
await pubsub.close()
await producer.stop()
loop.add_signal_handler(signal.SIGTERM, lambda: asyncio.create_task(graceful_shutdown()))Kafka 컨슈머들
같은 토픽을 서로 다른 컨슈머 그룹으로 읽는다.
- push-worker → 오프라인 유저에게 FCM/APNs
- index-worker → 검색 인덱스
메시지 영속화는 여기 없다. ws 파드가 채번과 같은 트랜잭션에서 이미 messages에 넣었기 때문이다. 별도 persist-worker를 두면 컨슈머 랙만큼 DB가 뒤처지고, 그 랙 구간이 그대로 복구 경로의 구멍이 된다. 클라이언트가 "seq 100 이후"를 요청했는데 DB엔 95까지만 있고 실시간으로는 103이 들어오는 상황이다. 채번과 저장을 붙여두면 이 문제가 원천적으로 없다.
push-worker의 수신자. 페이로드에는 발신자만 있으므로 방 멤버 목록을 따로 조회해야 한다.
members = await db.fetch("SELECT user_id FROM room_members WHERE room_id = $1", room_id)user:{uid}:online 조회로 오프라인만 거르는 방식은 판정과 실제 수신 사이에 레이스가 있다. 온라인이라 판정된 직후 끊긴 유저는 푸시도 실시간도 못 받는다. 미읽음 기반 지연 푸시가 안전하다. 몇 초 뒤에 다시 봐서 해당 seq가 여전히 안 읽혔으면 그때 보낸다.
이 구조가 성립하려면 필요한 것들
seq 채번이 사실상 전체를 지탱한다. 연속된 정수라서 클라이언트가 seq 7 다음에 9가 왔다 = 8을 놓쳤다를 즉시 안다. Snowflake나 ULID 같은 시간 기반 ID는 조율 지점이 없어 무한히 확장되지만 이 갭 감지가 원리적으로 불가능하다. Discord와 Slack은 그쪽을 골랐고, Telegram(pts)과 Matrix(stream_ordering)는 이쪽을 골랐다. 이 문서의 복구 경로 전체가 갭 감지 위에 얹혀 있으므로 후자를 택한다.
남은 유실 구간은 하나뿐이다. Redis pub/sub은 at-most-once라 재연결 틈에 발행된 메시지는 그냥 사라진다. 실시간 경로와 저장 경로는 이제 독립이 아니라 순차 종속(커밋 → ack → publish)이므로 순서가 갈라지지는 않는다. 클라이언트가 "내 마지막 seq 이후 주세요"로 DB에서 갭을 메우는 복구 경로가 반드시 짝으로 있어야 한다. DB가 항상 최신이므로 이 조회는 언제 해도 정확하다.
clientMsgId 멱등성. 재전송 시 같은 ID를 보낸다. 첫 전송이 성공했는데 응답만 유실된 케이스가 흔하다. 강제 지점은 UNIQUE(room_id, client_msg_id)이고, 채번 전에 기존 행을 조회해서 이미 있으면 그 seq를 그대로 돌려준다. 조회를 건너뛰면 중복 저장은 막히지만 seq는 이미 소비돼서 영구 결번이 남는다.
핫 파티션. 초대형 방 하나가 파티션 하나를 밀어붙이면 별도 토픽 분리나 키 살팅으로 대응. 같은 방은 room_seq 행 잠금도 직렬화 지점이 되므로, 초당 수백 건이 꽂히는 라이브 방송형 방은 아예 별도 경로로 빼는 편이 낫다.
클라이언트 재접속 로직. 서버가 after_ms를 못 보낸 경우(비정상 종료)에도 클라이언트가 알아서 지터 섞인 지수 백오프해야 함. 고정 지연이면 결국 동기화돼서 파도가 만들어진다.