From b17555a543d2b255d08688675844aeaa4458173a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=ED=99=8D=EC=84=B1=ED=9C=98?= Date: Mon, 3 Aug 2026 15:16:05 +0900 Subject: [PATCH 1/6] feat: add Redis pub/sub messaging infra for websocket --- .../realtime/config/RedisPubSubConfig.java | 40 +++++++++++++++++++ .../pubsub/RedisMessagePublisher.java | 35 ++++++++++++++++ .../pubsub/RedisMessageSubscriber.java | 32 +++++++++++++++ 3 files changed, 107 insertions(+) create mode 100644 momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java create mode 100644 momogo-realtime/src/main/java/com/momogo/realtime/websocket/pubsub/RedisMessagePublisher.java create mode 100644 momogo-realtime/src/main/java/com/momogo/realtime/websocket/pubsub/RedisMessageSubscriber.java diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java b/momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java new file mode 100644 index 0000000..6525077 --- /dev/null +++ b/momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java @@ -0,0 +1,40 @@ +package com.momogo.realtime.config; + +import lombok.RequiredArgsConstructor; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.data.redis.connection.RedisConnectionFactory; +import org.springframework.data.redis.listener.ChannelTopic; +import org.springframework.data.redis.listener.RedisMessageListenerContainer; +import org.springframework.data.redis.listener.adapter.MessageListenerAdapter; +import com.momogo.realtime.websocket.pubsub.RedisMessageSubscriber; + +@Configuration +@RequiredArgsConstructor +public class RedisPubSubConfig { + + // 실시간 시험방 통신에 사용할 Redis 채널 토픽 정의 + @Bean + public ChannelTopic roomTopic() { + return new ChannelTopic("room-realtime-channel"); + } + + // Redis Message를 수신할 어댑터 설정 (subscriber의 handleMessage 메소드를 호출) + @Bean + public MessageListenerAdapter listenerAdapter(RedisMessageSubscriber subscriber) { + return new MessageListenerAdapter(subscriber, "handleMessage"); + } + + @Bean + public RedisMessageListenerContainer redisMessageListenerContainer( + RedisConnectionFactory connectionFactory, + MessageListenerAdapter listenerAdapter, + ChannelTopic roomTopic + ) { + RedisMessageListenerContainer container = new RedisMessageListenerContainer(); + container.setConnectionFactory(connectionFactory); + // 토픽과 리스너 어댑터를 바인딩하여 메시지 수신 활성화 + container.addMessageListener(listenerAdapter, roomTopic); + return container; + } +} diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/pubsub/RedisMessagePublisher.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/pubsub/RedisMessagePublisher.java new file mode 100644 index 0000000..9551237 --- /dev/null +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/pubsub/RedisMessagePublisher.java @@ -0,0 +1,35 @@ +package com.momogo.realtime.websocket.pubsub; + +import com.fasterxml.jackson.databind.ObjectMapper; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.data.redis.core.RedisTemplate; +import org.springframework.data.redis.listener.ChannelTopic; +import org.springframework.stereotype.Component; + +@Slf4j +@Component +@RequiredArgsConstructor +public class RedisMessagePublisher { + + private final RedisTemplate redisTemplate; + private final ChannelTopic roomTopic; + private final ObjectMapper objectMapper; + + /** + * 메시지를 JSON으로 직렬화하여 Redis 채널로 Publish 합니다. + * @param messagePayload + */ + public void publish(Object messagePayload) { + try { + String jsonMessage = objectMapper.writeValueAsString(messagePayload); + log.info("[Redis Publisher] Redis 채널({})로 메시지 발행: {}", roomTopic.getTopic(), jsonMessage); + + redisTemplate.convertAndSend(roomTopic.getTopic(), jsonMessage); + + } catch (Exception e) { + log.error("[Redis Publisher] 메시지 발행 실패", e); + } + } + +} diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/pubsub/RedisMessageSubscriber.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/pubsub/RedisMessageSubscriber.java new file mode 100644 index 0000000..48d1351 --- /dev/null +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/pubsub/RedisMessageSubscriber.java @@ -0,0 +1,32 @@ +package com.momogo.realtime.websocket.pubsub; + +import com.fasterxml.jackson.databind.ObjectMapper; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.messaging.simp.SimpMessagingTemplate; +import org.springframework.stereotype.Component; + +@Slf4j +@Component +@RequiredArgsConstructor +public class RedisMessageSubscriber { + + private final ObjectMapper objectMapper; + private final SimpMessagingTemplate messagingTemplate; + + /** + * RedisPubSubConfig의 MessageListenerAdapter에 의해 호출되는 메소드 + * Redis에서 발행된 메시지를 받아 웹소켓 구독자(/sub/...)들에게 전파합니다. + * @param messageJson + */ + public void handleMessage(String messageJson) { + try { + log.info("[Redis Subscriber] 수신된 Pub/Sub 메시지: {}", messageJson); + + // TODO: JSON 메시지를 DTO로 직렬화해서 목적지(/sub/rooms/{roomId})로 웹소켓 브로드캐스팅 + } catch (Exception e) { + log.error("[Redis Subscriber] 메시지 처리 중 오류 발생", e); + } + } + +} From 020d9031bba3c60d080e97c095da9dfebd68c349 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=ED=99=8D=EC=84=B1=ED=9C=98?= Date: Mon, 3 Aug 2026 16:02:14 +0900 Subject: [PATCH 2/6] feat: impl STOMP websocket Redis pub/sub relay with DTO records and constants --- .../realtime/config/RedisPubSubConfig.java | 2 +- .../constant/WebSocketConstants.java | 10 +++++ .../controller/RoomRealtimeController.java | 31 ++++++++++++++ .../dto/request/RealtimeMessageRequest.java | 22 ++++++++++ .../dto/response/RealtimeMessageResponse.java | 42 +++++++++++++++++++ .../dto/type/RoomRealtimeStatus.java | 19 +++++++++ .../RedisMessagePublisher.java | 2 +- .../RedisMessageSubscriber.java | 11 ++++- 8 files changed, 135 insertions(+), 4 deletions(-) create mode 100644 momogo-realtime/src/main/java/com/momogo/realtime/websocket/constant/WebSocketConstants.java create mode 100644 momogo-realtime/src/main/java/com/momogo/realtime/websocket/controller/RoomRealtimeController.java create mode 100644 momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/request/RealtimeMessageRequest.java create mode 100644 momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/response/RealtimeMessageResponse.java create mode 100644 momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/type/RoomRealtimeStatus.java rename momogo-realtime/src/main/java/com/momogo/realtime/websocket/{pubsub => redis}/RedisMessagePublisher.java (95%) rename momogo-realtime/src/main/java/com/momogo/realtime/websocket/{pubsub => redis}/RedisMessageSubscriber.java (62%) diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java b/momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java index 6525077..2b00795 100644 --- a/momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java +++ b/momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java @@ -7,7 +7,7 @@ import org.springframework.data.redis.listener.ChannelTopic; import org.springframework.data.redis.listener.RedisMessageListenerContainer; import org.springframework.data.redis.listener.adapter.MessageListenerAdapter; -import com.momogo.realtime.websocket.pubsub.RedisMessageSubscriber; +import com.momogo.realtime.websocket.redis.RedisMessageSubscriber; @Configuration @RequiredArgsConstructor diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/constant/WebSocketConstants.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/constant/WebSocketConstants.java new file mode 100644 index 0000000..fa22331 --- /dev/null +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/constant/WebSocketConstants.java @@ -0,0 +1,10 @@ +package com.momogo.realtime.websocket.constant; + +public final class WebSocketConstants { + + public static final String SUB_ROOM_PREFIX = "/sub/rooms/"; + + private WebSocketConstants() { + + } +} diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/controller/RoomRealtimeController.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/controller/RoomRealtimeController.java new file mode 100644 index 0000000..1b8b274 --- /dev/null +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/controller/RoomRealtimeController.java @@ -0,0 +1,31 @@ +package com.momogo.realtime.websocket.controller; + +import com.momogo.realtime.websocket.dto.request.RealtimeMessageRequest; +import com.momogo.realtime.websocket.dto.response.RealtimeMessageResponse; +import com.momogo.realtime.websocket.redis.RedisMessagePublisher; +import java.util.UUID; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.messaging.handler.annotation.DestinationVariable; +import org.springframework.messaging.handler.annotation.MessageMapping; +import org.springframework.stereotype.Controller; + +@Slf4j +@Controller +@RequiredArgsConstructor +public class RoomRealtimeController { + + private final RedisMessagePublisher redisMessagePublisher; + + @MessageMapping("/room/{roomId}/send") + public void sendMessage( + @DestinationVariable("roomId")UUID roomId, + RealtimeMessageRequest request + ) { + log.info("[RoomRealtimeController] 시험방({}) 상태 변경 수신 - userId: {}, status: {}", + roomId, request.userId(), request.status()); + + RealtimeMessageResponse response = RealtimeMessageResponse.from(request); + redisMessagePublisher.publish(response); + } +} diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/request/RealtimeMessageRequest.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/request/RealtimeMessageRequest.java new file mode 100644 index 0000000..e591608 --- /dev/null +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/request/RealtimeMessageRequest.java @@ -0,0 +1,22 @@ +package com.momogo.realtime.websocket.dto.request; + +import com.momogo.realtime.websocket.dto.type.RoomRealtimeStatus; +import java.util.UUID; + +/** + * 웹소켓 실시간 응시 상태 변경 요청 DTO + * @param roomId 시험방 ID + * @param userId 수험생 ID + * @param status 응시 상태 + * @param solvedCount 푼 문제 개수 + * @param details 기타 실시간 상세 정보 + */ +public record RealtimeMessageRequest( + UUID roomId, + UUID userId, + RoomRealtimeStatus status, + int solvedCount, + Object details +) { + +} diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/response/RealtimeMessageResponse.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/response/RealtimeMessageResponse.java new file mode 100644 index 0000000..2256be8 --- /dev/null +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/response/RealtimeMessageResponse.java @@ -0,0 +1,42 @@ +package com.momogo.realtime.websocket.dto.response; + +import com.momogo.realtime.websocket.dto.request.RealtimeMessageRequest; +import com.momogo.realtime.websocket.dto.type.RoomRealtimeStatus; +import java.time.LocalDateTime; +import java.time.OffsetDateTime; +import java.util.UUID; + +/** + * 웹소켓 실시간 응시 상태 브로드캐스팅 응답 DTO + * @param roomId 시험방 ID + * @param userId 수험생 ID + * @param status 응시 상태 + * @param solvedCount 푼 문제 개수 + * @param details 기타 실시간 상세 정보 + * @param timestamp 응답 시간 + */ +public record RealtimeMessageResponse( + UUID roomId, + UUID userId, + RoomRealtimeStatus status, + int solvedCount, + Object details, + OffsetDateTime timestamp +) { + + /** + * Request DTO로부터 Response DTO를 정적 팩토리 메서드로 생성 + * @param request + * @return + */ + public static RealtimeMessageResponse from(RealtimeMessageRequest request) { + return new RealtimeMessageResponse( + request.roomId(), + request.userId(), + request.status(), + request.solvedCount(), + request.details(), + OffsetDateTime.now() + ); + } +} diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/type/RoomRealtimeStatus.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/type/RoomRealtimeStatus.java new file mode 100644 index 0000000..0c14584 --- /dev/null +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/type/RoomRealtimeStatus.java @@ -0,0 +1,19 @@ +package com.momogo.realtime.websocket.dto.type; + +import lombok.Getter; +import lombok.RequiredArgsConstructor; + +/** + * 실시간 시험방 응시 상태 구분용 ENUM + */ +@Getter +@RequiredArgsConstructor +public enum RoomRealtimeStatus { + + ENTER("시험방 입장"), + PROGRESS("문제 풀이 진행 중"), + SUBMIT("답안 제출 완료"), + LEAVE("시험방 퇴장"); + + private final String description; +} diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/pubsub/RedisMessagePublisher.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessagePublisher.java similarity index 95% rename from momogo-realtime/src/main/java/com/momogo/realtime/websocket/pubsub/RedisMessagePublisher.java rename to momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessagePublisher.java index 9551237..c8a79a6 100644 --- a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/pubsub/RedisMessagePublisher.java +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessagePublisher.java @@ -1,4 +1,4 @@ -package com.momogo.realtime.websocket.pubsub; +package com.momogo.realtime.websocket.redis; import com.fasterxml.jackson.databind.ObjectMapper; import lombok.RequiredArgsConstructor; diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/pubsub/RedisMessageSubscriber.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessageSubscriber.java similarity index 62% rename from momogo-realtime/src/main/java/com/momogo/realtime/websocket/pubsub/RedisMessageSubscriber.java rename to momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessageSubscriber.java index 48d1351..909a82a 100644 --- a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/pubsub/RedisMessageSubscriber.java +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessageSubscriber.java @@ -1,6 +1,8 @@ -package com.momogo.realtime.websocket.pubsub; +package com.momogo.realtime.websocket.redis; import com.fasterxml.jackson.databind.ObjectMapper; +import com.momogo.realtime.websocket.constant.WebSocketConstants; +import com.momogo.realtime.websocket.dto.response.RealtimeMessageResponse; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.messaging.simp.SimpMessagingTemplate; @@ -23,7 +25,12 @@ public void handleMessage(String messageJson) { try { log.info("[Redis Subscriber] 수신된 Pub/Sub 메시지: {}", messageJson); - // TODO: JSON 메시지를 DTO로 직렬화해서 목적지(/sub/rooms/{roomId})로 웹소켓 브로드캐스팅 + RealtimeMessageResponse response = objectMapper.readValue(messageJson, RealtimeMessageResponse.class); + String destination = WebSocketConstants.SUB_ROOM_PREFIX + response.roomId(); + + messagingTemplate.convertAndSend(destination, response); + log.info("[WebSocket Broadcast] 목적지({})로 상태 메시지 전파 완료", destination); + } catch (Exception e) { log.error("[Redis Subscriber] 메시지 처리 중 오류 발생", e); } From 6d06f9e7f892bc97b4946979f154b15063573ce6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=ED=99=8D=EC=84=B1=ED=9C=98?= Date: Mon, 3 Aug 2026 17:59:55 +0900 Subject: [PATCH 3/6] refactor: apply coderabbit security, validation, exception handling reviews for realtime websocket --- .../common/exception/RealtimeErrorCode.java | 27 +++++++++++++++++++ .../controller/RoomRealtimeController.java | 11 ++++++-- .../dto/request/RealtimeMessageRequest.java | 4 +++ .../dto/response/RealtimeMessageResponse.java | 6 ++--- .../redis/RedisMessagePublisher.java | 3 +++ 5 files changed, 46 insertions(+), 5 deletions(-) create mode 100644 momogo-core/src/main/java/com/momogo/core/common/exception/RealtimeErrorCode.java diff --git a/momogo-core/src/main/java/com/momogo/core/common/exception/RealtimeErrorCode.java b/momogo-core/src/main/java/com/momogo/core/common/exception/RealtimeErrorCode.java new file mode 100644 index 0000000..e66ec25 --- /dev/null +++ b/momogo-core/src/main/java/com/momogo/core/common/exception/RealtimeErrorCode.java @@ -0,0 +1,27 @@ +package com.momogo.core.common.exception; + +import lombok.Getter; +import lombok.RequiredArgsConstructor; +import org.springframework.http.HttpStatus; + +@Getter +@RequiredArgsConstructor +public enum RealtimeErrorCode implements ErrorCode { + + REDIS_PUBLISH_FAILED(8001, "REDIS_PUBLISH_FAILED", HttpStatus.SERVICE_UNAVAILABLE, "실시간 메시지 발행 중 오류가 발생했습니다."); + + private final int numeric; + private final String errorKey; + private final HttpStatus httpStatus; + private final String message; + + @Override + public String getDomain() { + return "REALTIME"; + } + + @Override + public String getCode() { + return getDomain() + "-" + getErrorKey(); + } +} diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/controller/RoomRealtimeController.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/controller/RoomRealtimeController.java index 1b8b274..3025639 100644 --- a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/controller/RoomRealtimeController.java +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/controller/RoomRealtimeController.java @@ -3,11 +3,14 @@ import com.momogo.realtime.websocket.dto.request.RealtimeMessageRequest; import com.momogo.realtime.websocket.dto.response.RealtimeMessageResponse; import com.momogo.realtime.websocket.redis.RedisMessagePublisher; +import jakarta.validation.Valid; +import java.security.Principal; import java.util.UUID; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.messaging.handler.annotation.DestinationVariable; import org.springframework.messaging.handler.annotation.MessageMapping; +import org.springframework.messaging.handler.annotation.Payload; import org.springframework.stereotype.Controller; @Slf4j @@ -20,12 +23,16 @@ public class RoomRealtimeController { @MessageMapping("/room/{roomId}/send") public void sendMessage( @DestinationVariable("roomId")UUID roomId, - RealtimeMessageRequest request + Principal principal, + @Payload @Valid RealtimeMessageRequest request ) { + + UUID authenticatedUserId = UUID.fromString(principal.getName()); + log.info("[RoomRealtimeController] 시험방({}) 상태 변경 수신 - userId: {}, status: {}", roomId, request.userId(), request.status()); - RealtimeMessageResponse response = RealtimeMessageResponse.from(request); + RealtimeMessageResponse response = RealtimeMessageResponse.of(roomId, authenticatedUserId, request); redisMessagePublisher.publish(response); } } diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/request/RealtimeMessageRequest.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/request/RealtimeMessageRequest.java index e591608..4f1a883 100644 --- a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/request/RealtimeMessageRequest.java +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/request/RealtimeMessageRequest.java @@ -1,6 +1,8 @@ package com.momogo.realtime.websocket.dto.request; import com.momogo.realtime.websocket.dto.type.RoomRealtimeStatus; +import jakarta.validation.constraints.Min; +import jakarta.validation.constraints.NotNull; import java.util.UUID; /** @@ -14,7 +16,9 @@ public record RealtimeMessageRequest( UUID roomId, UUID userId, + @NotNull(message = "응시 상태(status)는 필수 입력값입니다.") RoomRealtimeStatus status, + @Min(value = 0, message = "푼 문제 개수(solvedCount)는 0 이상이어야 합니다.") int solvedCount, Object details ) { diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/response/RealtimeMessageResponse.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/response/RealtimeMessageResponse.java index 2256be8..5499031 100644 --- a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/response/RealtimeMessageResponse.java +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/response/RealtimeMessageResponse.java @@ -29,10 +29,10 @@ public record RealtimeMessageResponse( * @param request * @return */ - public static RealtimeMessageResponse from(RealtimeMessageRequest request) { + public static RealtimeMessageResponse of(UUID roomId, UUID userId, RealtimeMessageRequest request) { return new RealtimeMessageResponse( - request.roomId(), - request.userId(), + roomId, + userId, request.status(), request.solvedCount(), request.details(), diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessagePublisher.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessagePublisher.java index c8a79a6..f663482 100644 --- a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessagePublisher.java +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessagePublisher.java @@ -1,6 +1,8 @@ package com.momogo.realtime.websocket.redis; import com.fasterxml.jackson.databind.ObjectMapper; +import com.momogo.core.common.exception.BusinessException; +import com.momogo.core.common.exception.RealtimeErrorCode; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.data.redis.core.RedisTemplate; @@ -29,6 +31,7 @@ public void publish(Object messagePayload) { } catch (Exception e) { log.error("[Redis Publisher] 메시지 발행 실패", e); + throw new BusinessException(RealtimeErrorCode.REDIS_PUBLISH_FAILED); } } From 258b9e5f815c6238d5fd45e565892fa938eac41b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=ED=99=8D=EC=84=B1=ED=9C=98?= Date: Mon, 3 Aug 2026 18:14:41 +0900 Subject: [PATCH 4/6] refactor: separate json serialization and redis publish exceptions in RedisMessagerPublisher --- .../common/exception/RealtimeErrorCode.java | 3 ++- .../websocket/redis/RedisMessagePublisher.java | 17 +++++++++++++---- 2 files changed, 15 insertions(+), 5 deletions(-) diff --git a/momogo-core/src/main/java/com/momogo/core/common/exception/RealtimeErrorCode.java b/momogo-core/src/main/java/com/momogo/core/common/exception/RealtimeErrorCode.java index e66ec25..698581d 100644 --- a/momogo-core/src/main/java/com/momogo/core/common/exception/RealtimeErrorCode.java +++ b/momogo-core/src/main/java/com/momogo/core/common/exception/RealtimeErrorCode.java @@ -8,7 +8,8 @@ @RequiredArgsConstructor public enum RealtimeErrorCode implements ErrorCode { - REDIS_PUBLISH_FAILED(8001, "REDIS_PUBLISH_FAILED", HttpStatus.SERVICE_UNAVAILABLE, "실시간 메시지 발행 중 오류가 발생했습니다."); + REDIS_PUBLISH_FAILED(8001, "REDIS_PUBLISH_FAILED", HttpStatus.SERVICE_UNAVAILABLE, "실시간 메시지 발행 중 오류가 발생했습니다."), + JSON_SERIALIZATION_FAILED(8002, "JSON_SERIALIZATION_FAILED", HttpStatus.INTERNAL_SERVER_ERROR, "실시간 메시지 직렬화 중 오류가 발생했습니다."); private final int numeric; private final String errorKey; diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessagePublisher.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessagePublisher.java index f663482..f123f94 100644 --- a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessagePublisher.java +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessagePublisher.java @@ -8,6 +8,7 @@ import org.springframework.data.redis.core.RedisTemplate; import org.springframework.data.redis.listener.ChannelTopic; import org.springframework.stereotype.Component; +import com.fasterxml.jackson.core.JsonProcessingException; @Slf4j @Component @@ -23,14 +24,22 @@ public class RedisMessagePublisher { * @param messagePayload */ public void publish(Object messagePayload) { + + String jsonMessage; + try { - String jsonMessage = objectMapper.writeValueAsString(messagePayload); - log.info("[Redis Publisher] Redis 채널({})로 메시지 발행: {}", roomTopic.getTopic(), jsonMessage); + jsonMessage = objectMapper.writeValueAsString(messagePayload); + } catch (JsonProcessingException e) { + log.error("[Redis Publisher] JSON 직렬화 실패 - payload: {}", messagePayload, e); + throw new BusinessException(RealtimeErrorCode.JSON_SERIALIZATION_FAILED); + } + // 2. Redis 메시지 발행 수행 (실패 시 503 REDIS_PUBLISH_FAILED) + try { + log.info("[Redis Publisher] Redis 채널({})로 메시지 발행: {}", roomTopic.getTopic(), jsonMessage); redisTemplate.convertAndSend(roomTopic.getTopic(), jsonMessage); - } catch (Exception e) { - log.error("[Redis Publisher] 메시지 발행 실패", e); + log.error("[Redis Publisher] Redis 메시지 전송 실패 - payload: {}", messagePayload, e); throw new BusinessException(RealtimeErrorCode.REDIS_PUBLISH_FAILED); } } From 273e2804bab2985a5e44c9980bdebf657180fc60 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=ED=99=8D=EC=84=B1=ED=9C=98?= Date: Tue, 4 Aug 2026 14:04:01 +0900 Subject: [PATCH 5/6] refactor: apply thread pool, error handler, dto refactor --- .../realtime/config/RedisPubSubConfig.java | 25 +++++++++++++++++++ .../controller/RoomRealtimeController.java | 4 +-- .../dto/request/RealtimeMessageRequest.java | 5 ---- .../redis/RedisMessageSubscriber.java | 2 +- 4 files changed, 28 insertions(+), 8 deletions(-) diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java b/momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java index 2b00795..91d49dc 100644 --- a/momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java +++ b/momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java @@ -1,6 +1,8 @@ package com.momogo.realtime.config; +import java.util.concurrent.Executor; import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.data.redis.connection.RedisConnectionFactory; @@ -8,7 +10,9 @@ import org.springframework.data.redis.listener.RedisMessageListenerContainer; import org.springframework.data.redis.listener.adapter.MessageListenerAdapter; import com.momogo.realtime.websocket.redis.RedisMessageSubscriber; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; +@Slf4j @Configuration @RequiredArgsConstructor public class RedisPubSubConfig { @@ -35,6 +39,27 @@ public RedisMessageListenerContainer redisMessageListenerContainer( container.setConnectionFactory(connectionFactory); // 토픽과 리스너 어댑터를 바인딩하여 메시지 수신 활성화 container.addMessageListener(listenerAdapter, roomTopic); + + // 커스텀 스레드 풀 적용 + Executor taskExecutor = createThreadPoolTaskExecutor(); + container.setTaskExecutor(taskExecutor); + container.setSubscriptionExecutor(taskExecutor); + + // 예외 에러 로깅 처리 + container.setErrorHandler(e -> + log.error("[Redis Pub/Sub Error] 비동기 메시지 수신/처리 중 오류 발생", e)); + return container; } + + private Executor createThreadPoolTaskExecutor() { + + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(10); + executor.setMaxPoolSize(50); + executor.setQueueCapacity(500); + executor.setThreadNamePrefix("redis-listener-"); + executor.initialize(); + return executor; + } } diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/controller/RoomRealtimeController.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/controller/RoomRealtimeController.java index 3025639..aa2e98b 100644 --- a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/controller/RoomRealtimeController.java +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/controller/RoomRealtimeController.java @@ -29,8 +29,8 @@ public void sendMessage( UUID authenticatedUserId = UUID.fromString(principal.getName()); - log.info("[RoomRealtimeController] 시험방({}) 상태 변경 수신 - userId: {}, status: {}", - roomId, request.userId(), request.status()); + log.info("[RoomRealtimeController] 시험방({}) 상태 변경 수신 - authenticatedUserId: {}, status: {}", + roomId, authenticatedUserId, request.status()); RealtimeMessageResponse response = RealtimeMessageResponse.of(roomId, authenticatedUserId, request); redisMessagePublisher.publish(response); diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/request/RealtimeMessageRequest.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/request/RealtimeMessageRequest.java index 4f1a883..388111f 100644 --- a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/request/RealtimeMessageRequest.java +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/dto/request/RealtimeMessageRequest.java @@ -3,19 +3,14 @@ import com.momogo.realtime.websocket.dto.type.RoomRealtimeStatus; import jakarta.validation.constraints.Min; import jakarta.validation.constraints.NotNull; -import java.util.UUID; /** * 웹소켓 실시간 응시 상태 변경 요청 DTO - * @param roomId 시험방 ID - * @param userId 수험생 ID * @param status 응시 상태 * @param solvedCount 푼 문제 개수 * @param details 기타 실시간 상세 정보 */ public record RealtimeMessageRequest( - UUID roomId, - UUID userId, @NotNull(message = "응시 상태(status)는 필수 입력값입니다.") RoomRealtimeStatus status, @Min(value = 0, message = "푼 문제 개수(solvedCount)는 0 이상이어야 합니다.") diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessageSubscriber.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessageSubscriber.java index 909a82a..54d4c0f 100644 --- a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessageSubscriber.java +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessageSubscriber.java @@ -21,7 +21,7 @@ public class RedisMessageSubscriber { * Redis에서 발행된 메시지를 받아 웹소켓 구독자(/sub/...)들에게 전파합니다. * @param messageJson */ - public void handleMessage(String messageJson) { + public void handleMessage(String messageJson) { try { log.info("[Redis Subscriber] 수신된 Pub/Sub 메시지: {}", messageJson); From 47d5a873f8b066a29c79843ccdde0af831cfc65f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=ED=99=8D=EC=84=B1=ED=9C=98?= Date: Tue, 4 Aug 2026 14:48:24 +0900 Subject: [PATCH 6/6] feat: check Redis Pub/Sub receiver count and log warning on zero subscribers --- .../realtime/config/RedisPubSubConfig.java | 61 ++++++++++++------- .../redis/RedisMessagePublisher.java | 21 ++++--- 2 files changed, 54 insertions(+), 28 deletions(-) diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java b/momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java index 91d49dc..37a6212 100644 --- a/momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java +++ b/momogo-realtime/src/main/java/com/momogo/realtime/config/RedisPubSubConfig.java @@ -1,15 +1,16 @@ package com.momogo.realtime.config; +import com.momogo.realtime.websocket.redis.RedisMessageSubscriber; import java.util.concurrent.Executor; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.data.redis.connection.RedisConnectionFactory; import org.springframework.data.redis.listener.ChannelTopic; import org.springframework.data.redis.listener.RedisMessageListenerContainer; import org.springframework.data.redis.listener.adapter.MessageListenerAdapter; -import com.momogo.realtime.websocket.redis.RedisMessageSubscriber; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; @Slf4j @@ -29,37 +30,55 @@ public MessageListenerAdapter listenerAdapter(RedisMessageSubscriber subscriber) return new MessageListenerAdapter(subscriber, "handleMessage"); } + /** + * 1. 개별 Pub/Sub 메시지 비동기 처리 전용 스레드 풀 (Spring Bean으로 등록하여 Graceful Shutdown 보장) + */ + @Bean(name = "redisTaskExecutor") + public ThreadPoolTaskExecutor redisTaskExecutor() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(10); + executor.setMaxPoolSize(50); + executor.setQueueCapacity(500); + executor.setThreadNamePrefix("redis-task-"); + executor.initialize(); + return executor; + } + + /** + * 2. Redis SUBSCRIBE 커넥션 롱폴링/유지 전용 독립 스레드 풀 (Spring Bean 등록) + */ + @Bean(name = "redisSubscriptionExecutor") + public ThreadPoolTaskExecutor redisSubscriptionExecutor() { + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(2); + executor.setMaxPoolSize(5); + executor.setQueueCapacity(10); + executor.setThreadNamePrefix("redis-sub-"); + executor.initialize(); + return executor; + } + @Bean public RedisMessageListenerContainer redisMessageListenerContainer( RedisConnectionFactory connectionFactory, MessageListenerAdapter listenerAdapter, - ChannelTopic roomTopic + ChannelTopic roomTopic, + @Qualifier("redisTaskExecutor") Executor redisTaskExecutor, + @Qualifier("redisSubscriptionExecutor") Executor redisSubscriptionExecutor ) { RedisMessageListenerContainer container = new RedisMessageListenerContainer(); container.setConnectionFactory(connectionFactory); - // 토픽과 리스너 어댑터를 바인딩하여 메시지 수신 활성화 container.addMessageListener(listenerAdapter, roomTopic); - // 커스텀 스레드 풀 적용 - Executor taskExecutor = createThreadPoolTaskExecutor(); - container.setTaskExecutor(taskExecutor); - container.setSubscriptionExecutor(taskExecutor); + // Spring Bean으로 생명주기가 안전하게 관리되는 두 스레드 풀 주입 + container.setTaskExecutor(redisTaskExecutor); + container.setSubscriptionExecutor(redisSubscriptionExecutor); - // 예외 에러 로깅 처리 + // 비동기 메시지 처리 중 예외 발생 시 에러 로깅 ErrorHandler container.setErrorHandler(e -> - log.error("[Redis Pub/Sub Error] 비동기 메시지 수신/처리 중 오류 발생", e)); + log.error("[Redis Pub/Sub Error] 비동기 메시지 수신/처리 중 오류 발생", e) + ); return container; } - - private Executor createThreadPoolTaskExecutor() { - - ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); - executor.setCorePoolSize(10); - executor.setMaxPoolSize(50); - executor.setQueueCapacity(500); - executor.setThreadNamePrefix("redis-listener-"); - executor.initialize(); - return executor; - } -} +} \ No newline at end of file diff --git a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessagePublisher.java b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessagePublisher.java index f123f94..e584ce1 100644 --- a/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessagePublisher.java +++ b/momogo-realtime/src/main/java/com/momogo/realtime/websocket/redis/RedisMessagePublisher.java @@ -1,21 +1,21 @@ package com.momogo.realtime.websocket.redis; +import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; import com.momogo.core.common.exception.BusinessException; import com.momogo.core.common.exception.RealtimeErrorCode; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.springframework.data.redis.core.RedisTemplate; +import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.data.redis.listener.ChannelTopic; import org.springframework.stereotype.Component; -import com.fasterxml.jackson.core.JsonProcessingException; @Slf4j @Component @RequiredArgsConstructor public class RedisMessagePublisher { - private final RedisTemplate redisTemplate; + private final StringRedisTemplate redisTemplate; private final ChannelTopic roomTopic; private final ObjectMapper objectMapper; @@ -30,16 +30,23 @@ public void publish(Object messagePayload) { try { jsonMessage = objectMapper.writeValueAsString(messagePayload); } catch (JsonProcessingException e) { - log.error("[Redis Publisher] JSON 직렬화 실패 - payload: {}", messagePayload, e); + log.error("[Redis Publisher] JSON 직렬화 실패", e); throw new BusinessException(RealtimeErrorCode.JSON_SERIALIZATION_FAILED); } // 2. Redis 메시지 발행 수행 (실패 시 503 REDIS_PUBLISH_FAILED) try { - log.info("[Redis Publisher] Redis 채널({})로 메시지 발행: {}", roomTopic.getTopic(), jsonMessage); - redisTemplate.convertAndSend(roomTopic.getTopic(), jsonMessage); + log.debug("[Redis Publisher] Redis 채널({})로 메시지 발행: {}", roomTopic.getTopic(), jsonMessage); + + Long receiversCount = redisTemplate.convertAndSend(roomTopic.getTopic(), jsonMessage); + + if (receiversCount == null || receiversCount == 0) { + log.warn("[Redis Publisher] 수신자 0명 - 채널({})로 발행되었으나 메시지를 수신한 구독자/인스턴스가 없습니다.", roomTopic.getTopic()); + } else { + log.info("[Redis Publisher] Redis 채널({}) 메시지 발행 성공 (수신 인스턴스: {}개)", roomTopic.getTopic(), receiversCount); + } } catch (Exception e) { - log.error("[Redis Publisher] Redis 메시지 전송 실패 - payload: {}", messagePayload, e); + log.error("[Redis Publisher] Redis 메시지 전송 실패 - topic: {}", roomTopic.getTopic(), e); throw new BusinessException(RealtimeErrorCode.REDIS_PUBLISH_FAILED); } }