Spring 숙련 (WebSocket)

KimGwangmin·2026년 9월 16일

세션 레지스트리

연결된 클라이언트의 WebSocketSession을 보관 및 조회하는 클래스

  • 앞서 배웠던 세션에서의 세션 저장소와 사실상 동일하다.
  • 거대한 Map

예제

RoomSessionRegistry

@Component
public class RoomSessionRegistry {
	// 방 번호(Long)별로 닉네임과 연결을 묶어 저장(Map<String, WebSocketSession>)
    // 같은 방에 연결된 대상들을 조회 가능
    // ConcurrentHashMap는 동시성 문제가 없는 HashMap, 대신 느림
    private final Map<Long, Map<String, WebSocketSession>> rooms = new ConcurrentHashMap<>();

	// 등록
    public boolean register(
            Long roomId,
            String nickname,
            WebSocketSession session
    ) {
        Map<String, WebSocketSession> sessionsByNickname = rooms.computeIfAbsent(roomId, id -> new ConcurrentHashMap<>());
        // putIfAbsent는 맵에 값이 없으면 넣으면서 null 반환, 있다면 해당 객체를 반환
        // 즉 이 조건은 새 등록이 성공한 경우 true를 반환함
        return sessionsByNickname.putIfAbsent(nickname, session) == null;
    }

	// 특정 방에 있는 대상을 조회
    public Collection<WebSocketSession> sessions(
            Long roomId
    ) {
        Map<String, WebSocketSession> sessionsByNickname = rooms.get(roomId);
        return sessionsByNickname == null ? List.of() : sessionsByNickname.values();
    }

    public boolean remove(
            Long roomId,
            String nickname,
            WebSocketSession session
    ) {
        Map<String, WebSocketSession> sessionsByNickname = rooms.get(roomId);
        return sessionsByNickname.remove(nickname, session);
    }
}

브로드캐스트

하나의 메시지를 여러 연결에 보내는 것

예제

@Slf4j
@Component
@RequiredArgsConstructor
public class RoomBroadcaster {
    private final RoomSessionRegistry registry;
    private final ObjectMapper objectMapper;

	// 특정 방에 속한 모든 세션에게 json 메시지 전송
    public void broadcastJson(
            Long roomId,
            String jsonMessage
    ) {
        for (WebSocketSession session : registry.sessions(roomId)) {
            sendJson(session, jsonMessage);
        }
    }

	// 특정 세션 하나에게 메시지 전송
    public void sendTo(
            WebSocketSession session,
            Object message
    ) {
        String jsonMessage = objectMapper.writeValueAsString(message);
        sendJson(session, jsonMessage);
    }

    private void sendJson(
            WebSocketSession session,
            String jsonMessage
    ) {
        try {
            synchronized (session) {
                if (session.isOpen()) {
                    TextMessage message = new TextMessage(jsonMessage);
                    session.sendMessage(message);
                }
            }
        } catch (Exception sendFailed) {
            log.warn("전송 실패: id={}", session.getId());
        }
    }
}

메시지 라우팅

메시지를 가장 알맞은 목적지로 전달하는 것

  • 메시지 종류마다 처리를 달리할 수 있음

메시지 종류 예시

  • {"type":"ping"}: 연결 상태 확인 요청
  • {"type":"chat","content":"안녕하세요"}: 채팅 요청

type 기반 메시지 라우팅

public interface WsMessageHandler {
    String type();

    void handle(
            WebSocketSession session,
            Long roomId,
            String nickname,
            JsonNode message // JSON 필드를 읽을 수 있는 객체
            // ObjectMapper가 받은 문자열을 JsonNode로 변환
            // JsonNode에선 각 필드값을 꺼낼 수 있다
    );
}

메시지 종류마다 처리 클래스를 만들고, 위 인터페이스를 구현하도록 하면 된다.

Heartbeat

연결이 응답하는지 주기적으로 확인하는 메시지
클라이언트가 일정 간격으로 ping을 보내면, 서버가 pong으로 응답하며 연결 상태를 확인
(반대로 서버가 heartbeat를 보내는 방식도 가능)

@Component
@RequiredArgsConstructor
public class PingWsHandler implements WsMessageHandler {
    private final RoomBroadcaster broadcaster;

    @Override
    public String type() {
        return "ping";
    }

    @Override
    public void handle(
            WebSocketSession session,
            Long roomId,
            String nickname,
            JsonNode message
    ) {
        PongResponse pong = new PongResponse();
        broadcaster.sendTo(session, pong); // 단일 응답
    }
}

메시지 수신 및 처리

  • ChatWebSocketHandler는 요청을 보고 적절한 핸들러를 찾아 넘겨준다.
    • 마치 HTTP 요청의 종류에 따라 적절한 컨트롤러를 찾아 요청을 넘겨주는 DispatcherServlet과 같은 역할을 한다.
    • 이때 바이트 배열인 메시지를 JSON으로 파싱해서 넘겨준다.
  • WsMessageRouter는 '요청을 보고 핸들러를 찾아주는' 기능을 담당한다.
@Slf4j
@Component
@RequiredArgsConstructor
public class ChatWebSocketHandler extends TextWebSocketHandler {
    private final RoomSessionRegistry registry;
    private final WsMessageRouter router;
    private final ObjectMapper objectMapper;

    @Override
    public void afterConnectionEstablished(
            WebSocketSession session
    ) throws Exception {
        // 연결 전에 Interceptor가 세션 attribute에 저장한 방 번호와 닉네임을 꺼냅니다.
        Long roomId = (Long) session.getAttributes().get(NicknameHandshakeInterceptor.ATTR_ROOM_ID);
        String nickname = (String) session.getAttributes().get(NicknameHandshakeInterceptor.ATTR_NICKNAME);
        // 같은 방에 닉네임이 이미 있으면 기존 연결을 유지하고 새 연결을 종료합니다.
        if (!registry.register(roomId, nickname, session)) {
            // 4002는 이 예제에서 닉네임 중복을 나타내도록 정한 종료 코드입니다.
            CloseStatus status = new CloseStatus(4002);
            session.close(status);
            return;
        }
        log.info("입장: room={} nickname={}", roomId, nickname);
    }

    @Override
    protected void handleTextMessage(
            WebSocketSession session,
            TextMessage message
    ) {
        Long roomId = (Long) session.getAttributes().get(NicknameHandshakeInterceptor.ATTR_ROOM_ID);
        String nickname = (String) session.getAttributes().get(NicknameHandshakeInterceptor.ATTR_NICKNAME);
        JsonNode parsedMessage;
        try {
            parsedMessage = objectMapper.readTree(message.getPayload());
        } catch (Exception notJson) {
            // JSON으로 읽을 수 없는 메시지는 처리하지 않습니다.
            return;
        }
        // 메시지를 처리할 핸들러 조회
        WsMessageHandler handler = router.find(parsedMessage);
        // type에 해당하는 핸들러가 없으면 이 메시지는 무시합니다.
        if (handler == null) {
            return;
        }
        handler.handle(session, roomId, nickname, parsedMessage);
    }

    @Override
    public void afterConnectionClosed(
            WebSocketSession session,
            CloseStatus status
    ) {
        Long roomId = (Long) session.getAttributes().get(NicknameHandshakeInterceptor.ATTR_ROOM_ID);
        String nickname = (String) session.getAttributes().get(NicknameHandshakeInterceptor.ATTR_NICKNAME);
        // 식별 정보가 없는 연결은 레지스트리에서 제거할 대상을 찾을 수 없습니다.
        if (roomId == null || nickname == null) {
            return;
        }
        // 현재 종료된 세션이 등록된 세션과 같을 때만 제거합니다.
        // 중복 닉네임으로 거절된 새 연결이 기존 연결을 지우는 것을 막습니다.
        if (!registry.remove(roomId, nickname, session)) {
            return;
        }
        log.info("퇴장: room={} nickname={} code={}", roomId, nickname, status.getCode());
    }
}
@Component
@RequiredArgsConstructor
public class WsMessageRouter {
    private final List<WsMessageHandler> handlers;

    public WsMessageHandler find(
            JsonNode message
    ) {
        String type = message.path("type").asString("");
        return handlers.stream()
                .filter(handler -> handler.type().equals(type))
                .findFirst()
                .orElse(null);
    }
}

라우터 클래스에서는 앞서 구현한 PingWsHandler 처럼 WsMessageHandler를 구현한 빈을 모아 리스트로 들고 있고, 조회 요청이 오면 리스트를 순회하며 타입이 일치하는 핸들러를 반환한다.
(PingWsHandler에 @Component가 붙어있는 것을 확인하자.)

실습

메시지 전송

  • WebSocket 핸들러는 메시지에서 필요한 값을 읽어 서비스에 전달한다.
    • HTTP 컨트롤러가 요청에서 값을 읽어 서비스 메서드를 호출하는 것과 비슷하다.

@Component
@RequiredArgsConstructor
public class ChatWsHandler implements WsMessageHandler {
    private final ChatService chatService;
    private final ChatRelay chatRelay;
    private final ObjectMapper objectMapper;

    @Override
    public String type() {
        return "chat";
    }

    @Override
    public void handle(
            WebSocketSession session,
            Long roomId,
            String nickname,
            JsonNode message
    ) {
        if (!message.path("content").isString()) {
            return;
        }
        String content = message.path("content").asString("");
        if (content.isBlank() || content.length() > 200) {
            return;
        }
        // 채팅을 DB에 저장
        ChatResponse savedChat = chatService.save(roomId, nickname, content);
        // savedChat을 JSON 문자열로 변환
        // savedChat은 타입, 발행자, 내용, 발행 시각이 담긴 DTO
        String jsonMessage = objectMapper.writeValueAsString(savedChat);
        // JSON 메시지를 Redis에 발행하여 같은 방의 연결에 전달
        // 해당 채널에 브로드캐스팅 (이건 DB를 거치지만 저장하는 건 아님, 브로드캐스팅만)
        chatRelay.publish(roomId, jsonMessage);
    }
}
@Slf4j
@Component
@RequiredArgsConstructor
public class ChatRelay implements MessageListener {
    private static final String CHANNEL_PREFIX = "chat:room:";
    private final StringRedisTemplate redisTemplate;
    private final RoomBroadcaster broadcaster;

    public void publish(
            Long roomId,
            String jsonMessage
    ) {
        redisTemplate.convertAndSend(CHANNEL_PREFIX + roomId, jsonMessage);
        log.info("발행: room={} {}", roomId, jsonMessage);
    }

    @Override
    public void onMessage(
            Message message,
            byte[] pattern
    ) {
        String channelName = new String(message.getChannel(), StandardCharsets.UTF_8);
        String jsonMessage = new String(message.getBody(), StandardCharsets.UTF_8);
        Long roomId = Long.valueOf(channelName.substring(CHANNEL_PREFIX.length()));
        log.info("수신: room={} {}", roomId, jsonMessage);
        broadcaster.broadcastJson(roomId, jsonMessage);
    }
}

실습

  • 좌측이 메시지를 보낸 서버, 우측은 별도의 서버이다.
  • 좌측 서버로 발행한 메시지를 우측 서버도 받는 것을 확인할 수 있다.

Redis가 필요한 이유

같은 방에 접속한 유저들이 모두 같은 서버에 연결되어 있으리란 보장은 없다!

1. 발행

  • 서버 A에 연결된 사용자가 방 1에 메시지를 전송
  • 서버 A는 Redis의 chat:room:1 채널로 메시지를 publish

2. 브로드캐스트

  • 방 1을 구독하고 있는 모든 서버(서버 A, 서버 B 등)의 ChatRelay.onMessage() 리스너가 해당 메시지를 수신

3. 로컬 발송

  • 각 서버는 자신의 로컬 메모리(RoomBroadcaster)를 뒤져 "현재 내 서버에 붙어 있는 방 1 세션들"에게만 session.sendMessage()를 호출

즉, 브로드캐스트가 두 단계로 진행된다.
1. DB에서 서버들에게
2. 서버 내에서 각 세션들에게

헷갈렸던 포인트 정리

Redis가 인메모리 DB라고 해서 '서버의 메모리를 사용한다'고 착각했었다. 그래서 여러 서버가 어떻게 하나의 DB를 공유하지? 라는 의문이 들었었다. 현재 로컬에서만 작업을 하고 있어 더 그런 오해가 생기는 것 같다. (물론 로컬이라고 스프링 서버의 메모리를 DB로 쓰는 것도 아니다)

인메모리 DB라는 것은, 스프링 서버와는 별개로 독립된 별도의 DB 프로세스가 돌아가는 서버 디바이스 안에서 디스크가 아닌 메모리를 사용한다는 의미이다. 따라서 스프링 서버의 입장에서는 구조적으로 전혀 달라지는 것이 없다. DB 서버 내부의 작동 방식이 다를 뿐이다.

일반적으로 디스크 기반 RDB에서는 서버-DB 간의 네트워크 통신 오버헤드보다 DB 내에서 디스크 I/O 오버헤드가 압도적으로 크다.
반면 메모리 기반 Redis에서는 전체 지연시간 중 DB 내에서의 I/O가 차지하는 비중이 0에 가깝고, 대부분이 네트워크 통신 시간이다.

0개의 댓글