연결된 클라이언트의
WebSocketSession을 보관 및 조회하는 클래스
MapRoomSessionRegistry
@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":"안녕하세요"}: 채팅 요청public interface WsMessageHandler {
String type();
void handle(
WebSocketSession session,
Long roomId,
String nickname,
JsonNode message // JSON 필드를 읽을 수 있는 객체
// ObjectMapper가 받은 문자열을 JsonNode로 변환
// JsonNode에선 각 필드값을 꺼낼 수 있다
);
}
메시지 종류마다 처리 클래스를 만들고, 위 인터페이스를 구현하도록 하면 된다.
연결이 응답하는지 주기적으로 확인하는 메시지
클라이언트가 일정 간격으로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는 요청을 보고 적절한 핸들러를 찾아 넘겨준다.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가 붙어있는 것을 확인하자.)


@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);
}
}


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

1. 발행
chat:room:1 채널로 메시지를 publish2. 브로드캐스트
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에 가깝고, 대부분이 네트워크 통신 시간이다.