WebSocket + STOMP + Redis Pub/Sub 기반 채팅 시스템 설계와 구현

E.NO·2026년 5월 4일

왜 이걸 만들었나

멘토-멘티 간 채팅 기능을 구현하면서 단순히 "기술 스택을 선택"하는 게 아니라,
각 선택에 트레이드오프가 있음을 직접 체감했다.
기술 선택의 이유, 실제 코드 구현, 그리고 남아 있는 한계까지 정리한다.


1. 실시간 통신 방식 비교 — 왜 WebSocket인가

Polling

클라이언트 → 서버 (1초마다) → "새 메시지 있어요?"
서버 → "없어요" (대부분의 경우)

메시지가 없어도 HTTP 요청이 계속 발생한다.
동시 접속자 수 × 요청 주기만큼 서버 부하가 쌓이고, 1초 폴링이면 최대 1초 지연이 생긴다.
채팅에서 1초 지연은 치명적이다.

SSE (Server-Sent Events)

  • 서버 → 클라이언트 단방향만 가능

  • 클라이언트가 메시지를 보내려면 별도 REST 요청이 필요

  • 채팅은 양방향이라 구조 자체가 맞지 않는다

    WebSocket ✅ ─

    HTTP Upgrade 핸드셰이크로 연결을 한 번 맺으면, 이후는 TCP 위에서 프레임 단위로 통신한다.
    불필요한 헤더 오버헤드 없이 저지연 양방향 통신이 가능해 채팅에 최적이다.

    방식방향연결 비용실시간성채팅 적합성
    Polling단방향매 요청마다 HTTP낮음
    SSE서버→클라이언트1회높음
    WebSocket양방향1회높음

    2. STOMP — raw WebSocket 위에 프로토콜을 얹은 이유

    WebSocket은 그 자체로는 "바이트를 주고받는 통로"일 뿐이다.
    메시지 라우팅, 구독, 브로드캐스트 같은 개념은 직접 구현해야 한다.

    STOMP(Simple Text Oriented Messaging Protocol)를 사용하면 이 부분을 프레임워크에 위임할 수 있다.

    WebSocket + STOMP 설정

    @Configuration
    @EnableWebSocketMessageBroker                                                                                                                                                                                           
    @RequiredArgsConstructor
    public class WebSocketConfig implements WebSocketMessageBrokerConfigurer {                                                                                                                                              
                                                                                                                                                                                                                          
        private final StompChannelInterceptor stompChannelInterceptor;
                                                                                                                                                                                                                            
        @Override
        public void registerStompEndpoints(StompEndpointRegistry registry) {                                                                                                                                                
            registry.addEndpoint("/ws-chat")  // 클라이언트 연결 엔드포인트                                                                                                                                               
                    .setAllowedOriginPatterns("*");                                                                                                                                                                         
        }                                                                                                                                                                                                                   
                                                                                                                                                                                                                            
        @Override                                                                                                                                                                                                           
        public void configureMessageBroker(MessageBrokerRegistry registry) {                                                                                                                                              
            registry.enableSimpleBroker("/topic", "/queue"); // 구독 prefix
            registry.setApplicationDestinationPrefixes("/app"); // 서버 처리 prefix                                                                                                                                         
        }                                                                                                                                                                                                                   
                                                                                                                                                                                                                            
        @Override                                                                                                                                                                                                           
        public void configureClientInboundChannel(ChannelRegistration registration) {                                                                                                                                     
            registration.interceptors(stompChannelInterceptor); // JWT 인증 인터셉터 등록
        }                                                                                                                                                                                                                   
    }
                                                                                                                                                                                                                      
  • /app/chat/{roomId} → 서버 컨트롤러(@MessageMapping)로 라우팅

  • /topic/chat/{roomId} → 해당 채팅방 구독자 전체에게 브로드캐스트

  • /queue/errors → 특정 사용자에게 에러 메시지 개별 전달


    3. WebSocket 인증 — ChannelInterceptor로 CONNECT 시점에 JWT 검증

    WebSocket은 HTTP와 달리 연결 이후 인증 메커니즘이 없다.
    Spring Security 필터 체인도 WebSocket 프레임에는 적용되지 않는다.

    해결: STOMP ChannelInterceptor

    @Slf4j                                                                                                                                                                                                                
    @Component
    @RequiredArgsConstructor
    public class StompChannelInterceptor implements ChannelInterceptor {
                                                                                                                                                                                                                            
        private final JwtUtil jwtUtil;
        private final UserDetailsService userDetailsService;                                                                                                                                                                
                                                                                                                                                                                                                            
        @Override
        public Message<?> preSend(Message<?> message, MessageChannel channel) {                                                                                                                                             
            StompHeaderAccessor accessor = MessageHeaderAccessor.getAccessor(message, StompHeaderAccessor.class);                                                                                                         
            if (accessor == null) return message;                                                                                                                                                                           
     
            if (StompCommand.CONNECT.equals(accessor.getCommand())) {                                                                                                                                                       
                String authHeader = accessor.getFirstNativeHeader("Authorization");                                                                                                                                         
     
                if (authHeader == null || !authHeader.startsWith(JwtUtil.BEARER_PREFIX)) {                                                                                                                                  
                    throw new ServiceErrorException(ChatExceptionEnum.ERR_WEBSOCKET_UNAUTHORIZED);                                                                                                                          
                }                                                                                                                                                                                                           
                                                                                                                                                                                                                            
                String token = jwtUtil.substringToken(authHeader);                                                                                                                                                          
                if (!jwtUtil.validateToken(token)) {                                                                                                                                                                        
                    throw new ServiceErrorException(ChatExceptionEnum.ERR_WEBSOCKET_UNAUTHORIZED);
                }                                                                                                                                                                                                           
                                                                                                                                                                                                                          
                String email = jwtUtil.extractEmail(token);
                UserDetails userDetails = userDetailsService.loadUserByUsername(email);                                                                                                                                     
                Long userId = ((UserDetailsImpl) userDetails).getUser().getId();
                                                                                                                                                                                                                            
                // userId를 Principal로 등록 → 컨트롤러에서 principal.getName()으로 접근                                                                                                                                  
                accessor.setUser(new UsernamePasswordAuthenticationToken(                                                                                                                                                   
                        userId.toString(), null, userDetails.getAuthorities()                                                                                                                                             
                ));                                                                                                                                                                                                         
                                                                                                                                                                                                                          
            } else if (StompCommand.SEND.equals(accessor.getCommand())                                                                                                                                                      
                    || StompCommand.SUBSCRIBE.equals(accessor.getCommand())) {                                                                                                                                              
                // CONNECT 없이 직접 SEND/SUBSCRIBE를 보내는 비인증 세션 차단
                if (accessor.getUser() == null) {                                                                                                                                                                           
                    throw new ServiceErrorException(ChatExceptionEnum.ERR_WEBSOCKET_UNAUTHORIZED);                                                                                                                        
                }                                                                                                                                                                                                           
            }                                                                                                                                                                                                             
                                                                                                                                                                                                                            
            return message;                                                                                                                                                                                               
        }
    }

    핵심 흐름:

  1. CONNECT 시점에 Authorization 헤더에서 JWT를 꺼내 검증

  2. 검증 성공 시 userId를 Principal로 세션에 저장

  3. 이후 SEND/SUBSCRIBE 시점에 accessor.getUser() == null 체크로 비정상 접근 차단

  4. 컨트롤러에서는 principal.getName()으로 userId를 꺼내 사용


    4. 메시지 송수신 — Controller와 Service

    WebSocket 컨트롤러

    @Slf4j                                                                                                                                                                                                                
    @Controller
    @RequiredArgsConstructor
    public class ChatController {
    
        private final ChatService chatService;
        private final ChatRedisPublisher chatRedisPublisher;
                                                                                                                                                                                                                            
        @MessageMapping("/chat/{roomId}")
        public void sendMessage(                                                                                                                                                                                            
                @DestinationVariable Long roomId,                                                                                                                                                                         
                @Valid @Payload ChatMessageRequest request,                                                                                                                                                                 
                Principal principal) {
            Long senderId = Long.parseLong(principal.getName()); // CONNECT 시 설정된 userId                                                                                                                                
                                                                                                                                                                                                                            
            ChatMessageResponse response = chatService.saveMessage(roomId, senderId, request);                                                                                                                              
            chatRedisPublisher.publish(roomId, ChatEventResponse.message(response));                                                                                                                                        
        }                                                                                                                                                                                                                   
                                                                                                                                                                                                                            
        @MessageExceptionHandler                                                                                                                                                                                            
        @SendToUser("/queue/errors")                                                                                                                                                                                      
        public String handleException(ServiceErrorException e) {                                                                                                                                                            
            return e.getMessage(); // 에러는 해당 유저의 /queue/errors로만 전달
        }                                                                                                                                                                                                                   
    }                                                                                                                                                                                                                       

    @MessageExceptionHandler + @SendToUser("/queue/errors")를 조합해
    에러 메시지가 다른 사용자에게 브로드캐스트되지 않도록 격리했다.

    메시지 저장 서비스

    @Transactional                                                                                                                                                                                                        
    public ChatMessageResponse saveMessage(Long roomId, Long senderId, ChatMessageRequest request) {
        if (request.messageType() == null || request.content() == null || request.content().isBlank()) {
            throw new ServiceErrorException(ChatExceptionEnum.ERR_INVALID_MESSAGE);                                                                                                                                         
        }
                                                                                                                                                                                                                            
        // FILE 타입이면 S3 버킷 URL 출처 검증 (링크 위·변조 방지)                                                                                                                                                        
        if (request.messageType() == MessageType.FILE) {                                                                                                                                                                    
            if (!isValidS3Url(request.fileUrl())) {
                throw new ServiceErrorException(ChatExceptionEnum.ERR_INVALID_MESSAGE);                                                                                                                                     
            }                                                                                                                                                                                                             
        }                                                                                                                                                                                                                   
                                                                                                                                                                                                                          
        ChatRoom room = getChatRoomOrThrow(roomId);
        validateActiveParticipant(room, senderId);                                                                                                                                                                          
                                                                                                                                                                                                                            
        ChatMessage message = (request.messageType() == MessageType.FILE)                                                                                                                                                   
                ? chatMessageRepository.save(                                                                                                                                                                               
                    ChatMessage.createWithFile(roomId, senderId, request.content(), request.messageType(), request.fileUrl()))                                                                                              
                : chatMessageRepository.save(                                                                                                                                                                               
                    ChatMessage.create(roomId, senderId, request.content(), request.messageType()));                                                                                                                        
                                                                                                                                                                                                                            
        return ChatMessageResponse.from(message);                                                                                                                                                                         
    }                                                                                                                                                                                                                       
     
    private boolean isValidS3Url(String fileUrl) {                                                                                                                                                                          
        String expectedPrefix = String.format("https://%s.s3.%s.amazonaws.com/chat/",                                                                                                                                     
                s3BucketName, s3Region);                                                                                                                                                                                    
        return fileUrl != null && fileUrl.startsWith(expectedPrefix);                                                                                                                                                       
    }                                                                                                                                                                                                                       

    DB 저장 → publish 순서를 반드시 지킨다.
    publish가 실패해도 DB에 메시지가 남아 있어 REST API로 재조회 및 복구가 가능하다.

    파일 메시지의 경우 S3 버킷 URL 패턴을 검증해서 외부 URL 삽입을 통한 위·변조를 방지했다.


    5. Redis Pub/Sub — 멀티 서버 브로드캐스트

    단일 서버라면 SimpMessagingTemplate.convertAndSend()만으로 충분하다.
    하지만 서버가 여러 대면 문제가 생긴다.

    [서버 A] 유저 1 접속 (WebSocket 연결)                                                                                                                                                                                   
    [서버 B] 유저 2 접속 (WebSocket 연결)                                                                                                                                                                                   
     
    유저 1 → 서버 A에 메시지 전송                                                                                                                                                                                           
    서버 A → 자신에게 연결된 클라이언트에게만 브로드캐스트                                                                                                                                                                  
    유저 2 (서버 B에 연결) → 메시지 못 받음 ❌                                                                                                                                                                               

    해결: Redis Pub/Sub으로 서버 간 메시지 중계

    서버 A → Redis publish("chat:room:1", message)                                                                                                                                                                        
    Redis  → 구독 중인 서버 A, B, C 모두에게 전달                                                                                                                                                                           
    각 서버 → /topic/chat/{roomId} 로 WebSocket 브로드캐스트                                                                                                                                                                
    클라이언트 수신 ✅                                                                                                                                                                                                      ─

    Publisher — Redis에 이벤트 발행

    @Slf4j                                                                                                                                                                                                                
    @Component
    @RequiredArgsConstructor
    public class ChatRedisPublisher {
    
        private static final int MAX_RETRY_ATTEMPTS = 3;
        private static final long RETRY_DELAY_MS = 100;                                                                                                                                                                     
     
        private final RedisTemplate<String, Object> redisTemplate;                                                                                                                                                          
        private final ObjectMapper objectMapper;                                                                                                                                                                            
     
        public void publish(Long roomId, ChatEventResponse event) {                                                                                                                                                         
            int attempt = 0;                                                                                                                                                                                              
            Exception lastException = null;
                                                                                                                                                                                                                            
            while (attempt < MAX_RETRY_ATTEMPTS) {
                try {                                                                                                                                                                                                       
                    byte[] channel = ("chat:room:" + roomId).getBytes(StandardCharsets.UTF_8);                                                                                                                            
                    byte[] message = objectMapper.writeValueAsBytes(event);                                                                                                                                                 
                    redisTemplate.execute((RedisCallback<Long>) conn -> conn.publish(channel, message));
                                                                                                                                                                                                                            
                    if (attempt > 0) {                                                                                                                                                                                    
                        log.info("[Redis] 이벤트 발행 성공 (재시도 {}회) roomId={}", attempt, roomId);                                                                                                                      
                    }                                                                                                                                                                                                       
                    return;
                                                                                                                                                                                                                            
                } catch (Exception e) {                                                                                                                                                                                   
                    lastException = e;
                    attempt++;
                    if (attempt < MAX_RETRY_ATTEMPTS) {
                        Thread.sleep(RETRY_DELAY_MS);                                                                                                                                                                       
                    }
                }                                                                                                                                                                                                           
            }                                                                                                                                                                                                             
    
            log.error("[Redis] 이벤트 발행 최종 실패 ({}회 재시도) roomId={}",
                    MAX_RETRY_ATTEMPTS, roomId, lastException);                                                                                                                                                             
        }
    }                                                                                                                                                                                                                       
  • 채널명 패턴: chat:room:{roomId}

  • 일시적인 Redis 장애에 대응해 3회, 100ms 간격 재시도 구현

  • 최종 실패 시 에러 로그 기록 (메시지는 DB에 이미 저장된 상태)

    Subscriber — Redis 수신 후 WebSocket 브로드캐스트

    @Slf4j                                                                                                                                                                                                                
    @Component
    @RequiredArgsConstructor
    public class ChatRedisSubscriber implements MessageListener {
    
        private final RedisMessageListenerContainer listenerContainer;
        private final SimpMessagingTemplate messagingTemplate;                                                                                                                                                              
        private final ObjectMapper objectMapper;
                                                                                                                                                                                                                            
        @PostConstruct                                                                                                                                                                                                    
        public void subscribe() {                                                                                                                                                                                           
            // 모든 채팅방 채널 구독 (패턴 매칭)
            listenerContainer.addMessageListener(this, new PatternTopic("chat:room:*"));                                                                                                                                    
        }                                                                                                                                                                                                                   
                                                                                                                                                                                                                            
        @Override                                                                                                                                                                                                           
        public void onMessage(Message message, byte[] pattern) {                                                                                                                                                          
            try {
                ChatEventResponse event = objectMapper.readValue(message.getBody(), ChatEventResponse.class);
                // 해당 채팅방을 구독 중인 모든 클라이언트에게 전달                                                                                                                                                         
                messagingTemplate.convertAndSend("/topic/chat/" + event.roomId(), event);
            } catch (Exception e) {                                                                                                                                                                                         
                log.error("[Redis] 메시지 처리 실패: {}", e.getMessage());                                                                                                                                                
            }                                                                                                                                                                                                               
        }                                                                                                                                                                                                                 
    }                                                                                                                                                                                                                       

    @PostConstruct로 애플리케이션 시작 시 chat:room:* 패턴을 구독하고,
    Redis에서 메시지를 받으면 즉시 해당 WebSocket 토픽으로 브로드캐스트한다.


    6. 이벤트 타입 통합 — ChatEventResponse

    클라이언트가 WebSocket으로 받는 이벤트는 단일 DTO로 통일했다.

    public record ChatEventResponse(                                                                                                                                                                                      
            String eventType,  // "MESSAGE" | "READ" | "LEAVE"
            Long userId,                                                                                                                                                                                                    
            Long roomId,
            ChatMessageResponse message  // MESSAGE 타입일 때만 존재, 나머지는 null                                                                                                                                         
    ) {                                                                                                                                                                                                                     
        public static ChatEventResponse message(ChatMessageResponse msg) {
            return new ChatEventResponse("MESSAGE", msg.senderId(), msg.roomId(), msg);                                                                                                                                     
        }                                                                                                                                                                                                                   
     
        public static ChatEventResponse read(Long roomId, Long readerId) {                                                                                                                                                  
            return new ChatEventResponse("READ", readerId, roomId, null);                                                                                                                                                   
        }
                                                                                                                                                                                                                            
        public static ChatEventResponse leave(Long roomId, Long userId) {                                                                                                                                                 
            return new ChatEventResponse("LEAVE", userId, roomId, null);
        }                                                                                                                                                                                                                   
    }

    메시지 수신, 읽음 처리, 나가기 이벤트를 하나의 DTO로 관리해서 클라이언트 처리 로직을 단순화했다.


    7. 동시성 처리 — 채팅방 중복 생성 방지

    같은 멘토십에 대해 동시에 채팅방 생성 요청이 오면 중복이 생길 수 있다.

    return chatRoomRepository.findByMentorshipId(mentorshipId)                                                                                                                                                            
            .orElseGet(() -> {                                                                                                                                                                                              
                try {                                                                                                                                                                                                       
                    // saveAndFlush로 즉시 DB 반영 → 동시 요청이 충돌을 감지할 수 있도록                                                                                                                                    
                    ChatRoom room = chatRoomRepository.saveAndFlush(                                                                                                                                                        
                            ChatRoom.create(mentorshipId, mentorUserId, menteeUserId)                                                                                                                                     
                    );                                                                                                                                                                                                      
                    return ChatRoomResponse.from(room);                                                                                                                                                                   
                } catch (DataIntegrityViolationException e) {                                                                                                                                                               
                    // 동시 요청으로 unique 제약 위반 시 이미 생성된 방 반환                                                                                                                                              
                    return ChatRoomResponse.from(                                                                                                                                                                           
                            chatRoomRepository.findByMentorshipId(mentorshipId)                                                                                                                                           
                                    .orElseThrow(...)                                                                                                                                                                       
                    );                                                                                                                                                                                                      
                }
            });                                                                                                                                                                                                             

    mentorshipId에 DB unique 제약을 걸어두고, 충돌 시 DataIntegrityViolationException을 잡아
    이미 생성된 방을 반환하는 방식으로 처리했다.


    8. 읽지 않은 메시지 수 — N+1 방지

    채팅방 목록 조회 시 각 방의 읽지 않은 메시지 수를 함께 보여줘야 한다.
    방마다 별도 쿼리를 날리면 N+1 문제가 생긴다.

    @Query("""                                                                                                                                                                                                              
        SELECT m.roomId, COUNT(m)                                                                                                                                                                                         
        FROM ChatMessage m
        WHERE m.roomId IN :roomIds
          AND m.senderId != :userId                                                                                                                                                                                         
          AND m.isRead = false
        GROUP BY m.roomId                                                                                                                                                                                                   
        """)                                                                                                                                                                                                              
    List<Object[]> countUnreadByRoomIds(@Param("roomIds") List<Long> roomIds, @Param("userId") Long userId);                                                                                                                

    방 ID 목록을 한 번에 넘겨서 IN 쿼리 하나로 모든 방의 읽지 않은 메시지 수를 가져온다.


    9. 전체 메시지 처리 흐름

    ① 클라이언트 STOMP SEND → /app/chat/{roomId}
    ② ChatController.sendMessage() 호출                                                                                                                                                                                     
    ③ ChatService.saveMessage() → DB 저장  ← 여기서 실패해도 유실 없음                                                                                                                                                      
    ④ ChatRedisPublisher.publish() → Redis에 이벤트 발행                                                                                                                                                                    
        └─ 실패 시 최대 3회 재시도 (100ms 간격)                                                                                                                                                                             
    ⑤ 모든 서버의 ChatRedisSubscriber.onMessage() 트리거                                                                                                                                                                    
    ⑥ SimpMessagingTemplate → /topic/chat/{roomId} 브로드캐스트                                                                                                                                                             
    ⑦ 클라이언트 수신                                                                                                                                                                                                       

    10. 한계점과 개선 방향

    ❗ Redis Pub/Sub은 메시지를 보관하지 않는다 ─

    Redis Pub/Sub은 fire-and-forget 구조다.
    구독 중이 아닌 서버나 클라이언트는 메시지를 받을 수 없다.
    서버가 잠깐 다운됐다가 올라오면 그 사이 발행된 메시지는 영원히 사라진다.

    개선 방향: Kafka 도입 → offset 기반으로 재처리 가능, 메시지를 디스크에 영속 보관

    ❗ 메시지 유실 가능성 ─

    Redis publish 직전에 서버가 다운되면 DB에는 저장됐지만 다른 사용자는 실시간으로 못 받는다.
    재시도와 DB 재조회 API로 보완하지만 완전한 at-least-once 보장은 아니다.

    개선 방향: Kafka의 acks + consumer group offset commit으로 정확한 전달 보장

    ❗ 메시지 순서 보장이 어렵다 ─

    멀티 서버 환경에서 네트워크 지연 차이로 메시지 순서가 뒤바뀔 수 있다.

    개선 방향: 클라이언트에서 createdAt 기준 정렬 / Kafka partition key로 순서 보장

    ❗ WebSocket 연결 끊김 ─

    모바일 환경 등 네트워크가 불안정한 경우 재연결 사이에 메시지를 놓친다.
    GET /{roomId}/messages REST API로 재조회해서 복구할 수 있지만, 클라이언트가 이를 직접 핸들링해야 한다.

    ❗ Redis Pub/Sub의 확장성 한계 ─

    fan-out 구조라 채널 수와 구독자 수가 늘어날수록 Redis 부하가 증가한다.

    개선 방향: Redis Stream 또는 Kafka로 전환


    11. 회고

    기술을 선택할 때 "이게 더 좋으니까"가 아니라
    "이 문제를 해결하기 위해 이 트레이드오프를 감수한다" 는 관점이 중요하다는 걸 느꼈다.

    WebSocket + STOMP는 실시간 양방향 통신 문제를 풀었고,
    Redis Pub/Sub은 멀티 서버 브로드캐스트 문제를 풀었다.
    DB 저장 선행 구조와 재시도 로직은 신뢰성을 보완했다.

    하지만 Redis Pub/Sub의 비영속성과 순서 보장 문제는 여전히 남아 있다.
    이걸 인식하고 어떻게 개선할 수 있는지 알고 있다는 것,
    그게 지금 단계에서 가장 중요한 성과라고 생각한다.


    한 줄 요약
    WebSocket + STOMP으로 실시간 양방향 채팅을 구현하고, Redis Pub/Sub으로 멀티 서버 브로드캐스트를 해결했다.
    DB 저장 선행 + publish 재시도 + REST 재조회 API로 메시지 유실을 보완했지만, 완전한 보장을 위해선 Kafka 전환이 필요하다.

0개의 댓글