멘토-멘티 간 채팅 기능을 구현하면서 단순히 "기술 스택을 선택"하는 게 아니라,
각 선택에 트레이드오프가 있음을 직접 체감했다.
기술 선택의 이유, 실제 코드 구현, 그리고 남아 있는 한계까지 정리한다.
클라이언트 → 서버 (1초마다) → "새 메시지 있어요?"
서버 → "없어요" (대부분의 경우)
메시지가 없어도 HTTP 요청이 계속 발생한다.
동시 접속자 수 × 요청 주기만큼 서버 부하가 쌓이고, 1초 폴링이면 최대 1초 지연이 생긴다.
채팅에서 1초 지연은 치명적이다.
서버 → 클라이언트 단방향만 가능
클라이언트가 메시지를 보내려면 별도 REST 요청이 필요
채팅은 양방향이라 구조 자체가 맞지 않는다
HTTP Upgrade 핸드셰이크로 연결을 한 번 맺으면, 이후는 TCP 위에서 프레임 단위로 통신한다.
불필요한 헤더 오버헤드 없이 저지연 양방향 통신이 가능해 채팅에 최적이다.
| 방식 | 방향 | 연결 비용 | 실시간성 | 채팅 적합성 |
|---|---|---|---|---|
| Polling | 단방향 | 매 요청마다 HTTP | 낮음 | ❌ |
| SSE | 서버→클라이언트 | 1회 | 높음 | ❌ |
| WebSocket | 양방향 | 1회 | 높음 | ✅ |
WebSocket은 그 자체로는 "바이트를 주고받는 통로"일 뿐이다.
메시지 라우팅, 구독, 브로드캐스트 같은 개념은 직접 구현해야 한다.
STOMP(Simple Text Oriented Messaging Protocol)를 사용하면 이 부분을 프레임워크에 위임할 수 있다.
@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 → 특정 사용자에게 에러 메시지 개별 전달
WebSocket은 HTTP와 달리 연결 이후 인증 메커니즘이 없다.
Spring Security 필터 체인도 WebSocket 프레임에는 적용되지 않는다.
@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;
}
}
핵심 흐름:
CONNECT 시점에 Authorization 헤더에서 JWT를 꺼내 검증
검증 성공 시 userId를 Principal로 세션에 저장
이후 SEND/SUBSCRIBE 시점에 accessor.getUser() == null 체크로 비정상 접근 차단
컨트롤러에서는 principal.getName()으로 userId를 꺼내 사용
@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 삽입을 통한 위·변조를 방지했다.
단일 서버라면 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 브로드캐스트
클라이언트 수신 ✅ ─
@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에 이미 저장된 상태)
@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 토픽으로 브로드캐스트한다.
클라이언트가 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로 관리해서 클라이언트 처리 로직을 단순화했다.
같은 멘토십에 대해 동시에 채팅방 생성 요청이 오면 중복이 생길 수 있다.
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을 잡아
이미 생성된 방을 반환하는 방식으로 처리했다.
채팅방 목록 조회 시 각 방의 읽지 않은 메시지 수를 함께 보여줘야 한다.
방마다 별도 쿼리를 날리면 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 쿼리 하나로 모든 방의 읽지 않은 메시지 수를 가져온다.
① 클라이언트 STOMP SEND → /app/chat/{roomId}
② ChatController.sendMessage() 호출
③ ChatService.saveMessage() → DB 저장 ← 여기서 실패해도 유실 없음
④ ChatRedisPublisher.publish() → Redis에 이벤트 발행
└─ 실패 시 최대 3회 재시도 (100ms 간격)
⑤ 모든 서버의 ChatRedisSubscriber.onMessage() 트리거
⑥ SimpMessagingTemplate → /topic/chat/{roomId} 브로드캐스트
⑦ 클라이언트 수신
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로 순서 보장
모바일 환경 등 네트워크가 불안정한 경우 재연결 사이에 메시지를 놓친다.
GET /{roomId}/messages REST API로 재조회해서 복구할 수 있지만, 클라이언트가 이를 직접 핸들링해야 한다.
fan-out 구조라 채널 수와 구독자 수가 늘어날수록 Redis 부하가 증가한다.
개선 방향: Redis Stream 또는 Kafka로 전환
기술을 선택할 때 "이게 더 좋으니까"가 아니라
"이 문제를 해결하기 위해 이 트레이드오프를 감수한다" 는 관점이 중요하다는 걸 느꼈다.
WebSocket + STOMP는 실시간 양방향 통신 문제를 풀었고,
Redis Pub/Sub은 멀티 서버 브로드캐스트 문제를 풀었다.
DB 저장 선행 구조와 재시도 로직은 신뢰성을 보완했다.
하지만 Redis Pub/Sub의 비영속성과 순서 보장 문제는 여전히 남아 있다.
이걸 인식하고 어떻게 개선할 수 있는지 알고 있다는 것,
그게 지금 단계에서 가장 중요한 성과라고 생각한다.
한 줄 요약
WebSocket + STOMP으로 실시간 양방향 채팅을 구현하고, Redis Pub/Sub으로 멀티 서버 브로드캐스트를 해결했다.
DB 저장 선행 + publish 재시도 + REST 재조회 API로 메시지 유실을 보완했지만, 완전한 보장을 위해선 Kafka 전환이 필요하다.