From 9da1aec24afc87157c428bbd54ffc949af71f7b6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=ED=99=8D=EC=84=B1=ED=9C=98?= Date: Tue, 11 Aug 2026 20:56:18 +0900 Subject: [PATCH 1/2] refactor(room): optimize answer submission with Kafka async processing and Redisson lock - Replace pessimistic DB lock (findByIdForUpdate) with Redisson distributed lock - Add ROOM_SUBMIT_EVENTS topic to KafkaTopics - Create RoomSubmitEventMessage DTO record with eventId idempotency key - Implement RoomSubmitKafkaProducer in momogo-core for async event publishing - Implement RoomSubmitKafkaConsumer in momogo-api with Redis deduplication - Refactor RoomServiceImpl.submitRoomAnswer to delegate DB writes to Kafka --- .../consumer/RoomSubmitKafkaConsumer.java | 93 +++++++++++++++++++ .../core/common/config/KafkaTopics.java | 1 + .../room/event/RoomSubmitEventMessage.java | 27 ++++++ .../producer/RoomSubmitKafkaProducer.java | 33 +++++++ .../domain/room/service/RoomServiceImpl.java | 93 ++++++++++--------- 5 files changed, 203 insertions(+), 44 deletions(-) create mode 100644 momogo-api/src/main/java/com/momogo/api/room/consumer/RoomSubmitKafkaConsumer.java create mode 100644 momogo-core/src/main/java/com/momogo/core/domain/room/event/RoomSubmitEventMessage.java create mode 100644 momogo-core/src/main/java/com/momogo/core/domain/room/kafka/producer/RoomSubmitKafkaProducer.java diff --git a/momogo-api/src/main/java/com/momogo/api/room/consumer/RoomSubmitKafkaConsumer.java b/momogo-api/src/main/java/com/momogo/api/room/consumer/RoomSubmitKafkaConsumer.java new file mode 100644 index 0000000..319c96a --- /dev/null +++ b/momogo-api/src/main/java/com/momogo/api/room/consumer/RoomSubmitKafkaConsumer.java @@ -0,0 +1,93 @@ +package com.momogo.api.room.consumer; + +import com.momogo.core.common.config.KafkaTopics; +import com.momogo.core.domain.room.entity.RoomProblem; +import com.momogo.core.domain.room.entity.RoomUser; +import com.momogo.core.domain.room.entity.RoomUserId; +import com.momogo.core.domain.room.entity.UserRoomAnswer; +import com.momogo.core.domain.room.event.RoomSubmitEventMessage; +import com.momogo.core.domain.room.repository.RoomProblemRepository; +import com.momogo.core.domain.room.repository.RoomUserRepository; +import com.momogo.core.domain.room.repository.UserRoomAnswerRepository; +import com.momogo.core.domain.user.entity.User; +import jakarta.persistence.EntityManager; +import java.time.Duration; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.function.Function; +import java.util.stream.Collectors; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.data.redis.core.RedisTemplate; +import org.springframework.kafka.annotation.KafkaListener; +import org.springframework.stereotype.Component; +import org.springframework.transaction.annotation.Transactional; +import org.springframework.transaction.support.TransactionSynchronization; +import org.springframework.transaction.support.TransactionSynchronizationManager; + +/* + * RoomServiceImpl이 발행한 답안 제출 메시지를 수신하여 실제 DB 저장을 비동기로 수행 + * Redis 멱등성 체크를 통해 중복 전달(At-least-once)로 인한 UQ 충돌 방지 + */ +@Slf4j +@Component +@RequiredArgsConstructor +public class RoomSubmitKafkaConsumer { + + private static final String DEDUP_KEY_PREFIX = "room:submit:processed:"; + private static final Duration DEDUP_TTL = Duration.ofHours(24); + + private final UserRoomAnswerRepository userRoomAnswerRepository; + private final RoomUserRepository roomUserRepository; + private final RoomProblemRepository roomProblemRepository; + private final RedisTemplate redisTemplate; + private final EntityManager entityManager; + + @KafkaListener(topics = KafkaTopics.ROOM_SUBMIT_EVENTS) + @Transactional + public void consume(RoomSubmitEventMessage message) { + log.info("[RoomSubmitKafkaConsumer] 카프카 답안 수신 - eventId: {}, userId: {}, roomId: {}", + message.eventId(), message.userId(), message.roomId()); + + // 1. Redis 멱등성 키 검증 (중복 처리 방지) + String dedupKey = DEDUP_KEY_PREFIX + message.eventId(); + if (Boolean.TRUE.equals(redisTemplate.hasKey(dedupKey))) { + log.warn("[RoomSubmitKafkaConsumer] 이미 처리된 답안 제출 이벤트, 중복 스킵 - eventId: {}", message.eventId()); + return; + } + + // 2. 추가 SELECT 없이 EntityManager 프록시 참조 생성 + User userProxy = entityManager.getReference(User.class, message.userId()); + + // 3. 문제 목록 조회 및 답안 일괄 매핑 + List problemIds = message.answers().stream() + .map(ans -> ans.roomProblemId()) + .toList(); + + Map problemsById = roomProblemRepository.findAllById(problemIds).stream() + .collect(Collectors.toMap(RoomProblem::getId, Function.identity())); + + List userRoomAnswers = message.answers().stream() + .map(ans -> UserRoomAnswer.of(userProxy, problemsById.get(ans.roomProblemId()), ans.userAnswer(), null)) + .toList(); + + // 4. DB 일괄 저장 + userRoomAnswerRepository.saveAll(userRoomAnswers); + + // 5. 응시 완료 상태 업데이트 + RoomUserId roomUserId = new RoomUserId(message.roomId(), message.userId()); + roomUserRepository.findById(roomUserId).ifPresent(RoomUser::attend); + + // 6. DB 커밋 확정 후에만 Redis 멱등성 처리 완료 키 갱신 + TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { + @Override + public void afterCommit() { + redisTemplate.opsForValue().set(dedupKey, "1", DEDUP_TTL); + } + }); + + log.info("[RoomSubmitKafkaConsumer] DB 저장 완료 - userId: {}, roomId: {}, 문항 수: {}", + message.userId(), message.roomId(), userRoomAnswers.size()); + } +} diff --git a/momogo-core/src/main/java/com/momogo/core/common/config/KafkaTopics.java b/momogo-core/src/main/java/com/momogo/core/common/config/KafkaTopics.java index 47f01fd..77c6033 100644 --- a/momogo-core/src/main/java/com/momogo/core/common/config/KafkaTopics.java +++ b/momogo-core/src/main/java/com/momogo/core/common/config/KafkaTopics.java @@ -8,6 +8,7 @@ public final class KafkaTopics { public static final String NOTIFICATION_EVENTS = "notification-events"; public static final String AI_GRADING_EVENTS = "ai-grading-events"; + public static final String ROOM_SUBMIT_EVENTS = "room-submit-events"; private KafkaTopics() { } diff --git a/momogo-core/src/main/java/com/momogo/core/domain/room/event/RoomSubmitEventMessage.java b/momogo-core/src/main/java/com/momogo/core/domain/room/event/RoomSubmitEventMessage.java new file mode 100644 index 0000000..87f60af --- /dev/null +++ b/momogo-core/src/main/java/com/momogo/core/domain/room/event/RoomSubmitEventMessage.java @@ -0,0 +1,27 @@ +package com.momogo.core.domain.room.event; + +import com.momogo.core.domain.room.dto.request.ProblemAnswerRequest; +import java.time.OffsetDateTime; +import java.util.List; +import java.util.Objects; +import java.util.UUID; + +public record RoomSubmitEventMessage( + UUID eventId, // 카프카 재전달 시 중복 처리 방지용 식별자 (Redis Dedup 키) + UUID userId, + UUID roomId, + List answers, + OffsetDateTime submittedAt +) { + + public RoomSubmitEventMessage { + Objects.requireNonNull(eventId, "eventId는 null일 수 없습니다"); + Objects.requireNonNull(userId, "userId는 null일 수 없습니다"); + Objects.requireNonNull(roomId, "roomId는 null일 수 없습니다"); + answers = List.copyOf(answers); + } + + public static RoomSubmitEventMessage of(UUID userId, UUID roomId, List answers) { + return new RoomSubmitEventMessage(UUID.randomUUID(), userId, roomId, answers, OffsetDateTime.now()); + } +} diff --git a/momogo-core/src/main/java/com/momogo/core/domain/room/kafka/producer/RoomSubmitKafkaProducer.java b/momogo-core/src/main/java/com/momogo/core/domain/room/kafka/producer/RoomSubmitKafkaProducer.java new file mode 100644 index 0000000..647fec9 --- /dev/null +++ b/momogo-core/src/main/java/com/momogo/core/domain/room/kafka/producer/RoomSubmitKafkaProducer.java @@ -0,0 +1,33 @@ +package com.momogo.core.domain.room.kafka.producer; + +import com.momogo.core.common.config.KafkaTopics; +import com.momogo.core.domain.room.event.RoomSubmitEventMessage; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.stereotype.Component; + +/* + * 답안 제출 이벤트를 room-submit-events 토픽에 비동기 발행하는 전담 클래스. + * 실제 DB 저장은 momogo-api의 RoomSubmitKafkaConsumer가 비동기로 처리 + */ +@Slf4j +@Component +@RequiredArgsConstructor +public class RoomSubmitKafkaProducer { + + private final KafkaTemplate kafkaTemplate; + + public void send(RoomSubmitEventMessage message) { + kafkaTemplate.send(KafkaTopics.ROOM_SUBMIT_EVENTS, message.roomId().toString(), message) + .whenComplete((result, ex) -> { + if (ex != null) { + log.error("[RoomSubmitKafkaProducer] 카프카 메시지 전송 실패 - eventId: {}, roomId: {}, topic: {}", + message.eventId(), message.roomId(), KafkaTopics.ROOM_SUBMIT_EVENTS, ex); + } else { + log.info("[RoomSubmitKafkaProducer] 카프카 메시지 전송 성공 - eventId: {}, offset: {}", + message.eventId(), result.getRecordMetadata().offset()); + } + }); + } +} diff --git a/momogo-core/src/main/java/com/momogo/core/domain/room/service/RoomServiceImpl.java b/momogo-core/src/main/java/com/momogo/core/domain/room/service/RoomServiceImpl.java index 36ad8c7..ef089de 100644 --- a/momogo-core/src/main/java/com/momogo/core/domain/room/service/RoomServiceImpl.java +++ b/momogo-core/src/main/java/com/momogo/core/domain/room/service/RoomServiceImpl.java @@ -63,6 +63,13 @@ import java.util.stream.Collectors; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import com.momogo.core.common.exception.AuthErrorCode; +import com.momogo.core.common.exception.GlobalErrorCode; +import com.momogo.core.domain.room.event.RoomSubmitEventMessage; +import com.momogo.core.domain.room.kafka.producer.RoomSubmitKafkaProducer; +import java.util.concurrent.TimeUnit; +import org.redisson.api.RLock; +import org.redisson.api.RedissonClient; import org.springframework.context.ApplicationEventPublisher; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; @@ -84,6 +91,10 @@ public class RoomServiceImpl implements RoomService{ private final RoomMapper roomMapper; private final ApplicationEventPublisher eventPublisher; private final AiGradingProducer aiGradingProducer; + private final RedissonClient redissonClient; + private final RoomSubmitKafkaProducer roomSubmitKafkaProducer; + + private static final String SUBMIT_LOCK_PREFIX = "lock:room:submit:"; // 검증된 타겟 객체들을 묶기 위한 private record private record ValidatedRoomTarget(Space space, List targetUsers) {} @@ -236,11 +247,10 @@ public List getRoomProblems(UUID userId, UUID roomId) { } @Override - @Transactional public void submitRoomAnswer(UUID userId, UUID roomId, RoomAnswerSubmitRequest request) { - log.info("[RoomService] 시험 답안 제출 시작 - userId: {}, roomId: {}", userId, roomId); + log.info("[RoomService] 시험 답안 제출 시도 - userId: {}, roomId: {}", userId, roomId); - // 요청 리스트 내 동일한 문제 ID가 중복 유입되는지 검증 + // 1. 요청 리스트 내 동일한 문제 ID가 중복 유입되는지 검증 long uniqueProblemCount = request.answers().stream() .map(ProblemAnswerRequest::roomProblemId) .distinct().count(); @@ -248,55 +258,50 @@ public void submitRoomAnswer(UUID userId, UUID roomId, RoomAnswerSubmitRequest r throw new BusinessException(RoomErrorCode.DUPLICATE_ANSWER_SUBMITTED); } - // 방 존재 검증 - Room room = roomRepository.findByIdForUpdate(roomId) - .orElseThrow(() -> new BusinessException(RoomErrorCode.ROOM_NOT_FOUND)); - - // 시간 검증 - 시험 시작 전 제출 시도 차단 - OffsetDateTime now = OffsetDateTime.now(); - if (now.isBefore(room.getTestStartAt())) { - throw new BusinessException(RoomErrorCode.INVALID_ACCESS_BEFORE_START); - } - - // 상태 검증 - 이미 최종적으로 끝난 시험이면 추가 제출 불가 - if (Boolean.TRUE.equals(room.getIsEnded())) { - throw new BusinessException(RoomErrorCode.ALREADY_ENDED); - } + // 2. Redisson 분산 락 획득 (Redisson Watchdog 활용을 위해 leaseTime 생략) + String lockKey = SUBMIT_LOCK_PREFIX + roomId + ":" + userId; + RLock lock = redissonClient.getLock(lockKey); - // 응시 대상 유저 자격 검증 - RoomUserId roomUserId = new RoomUserId(roomId, userId); - RoomUser roomUser = roomUserRepository.findByIdForUpdate(roomUserId) - .orElseThrow(() -> new BusinessException(RoomErrorCode.NOT_ROOM_PARTICIPANT)); + try { + if (lock.tryLock(3, TimeUnit.SECONDS)) { + try { + // 3. 비관적 락(findByIdForUpdate) 대신 락이 없는 일반 조회 사용 + Room room = findRoomOrThrow(roomId); - // 요청 유저 획득 - User user = userRepository.findById(userId) - .orElseThrow(() -> new BusinessException(SpaceErrorCode.SPACE_USER_NOT_FOUND)); + // 4. 시간 및 응시 자격 빠른 검증 + OffsetDateTime now = OffsetDateTime.now(); + if (now.isBefore(room.getTestStartAt())) { + throw new BusinessException(RoomErrorCode.INVALID_ACCESS_BEFORE_START); + } - List problemIds = request.answers().stream().map(ProblemAnswerRequest::roomProblemId).toList(); - Map problemsById = roomProblemRepository.findAllById(problemIds).stream() - .collect(Collectors.toMap(RoomProblem::getId, Function.identity())); + if (Boolean.TRUE.equals(room.getIsEnded())) { + throw new BusinessException(RoomErrorCode.ALREADY_ENDED); + } - if (Boolean.TRUE.equals(roomUser.getIsAttended())) { - throw new BusinessException(RoomErrorCode.ALREADY_ENDED); - } + RoomUserId roomUserId = new RoomUserId(roomId, userId); + RoomUser roomUser = roomUserRepository.findById(roomUserId) + .orElseThrow(() -> new BusinessException(RoomErrorCode.NOT_ROOM_PARTICIPANT)); - // 각 문제 답안 일괄 매핑, 저장 - List userRoomAnswers = request.answers().stream() - .map(ans -> { - RoomProblem problem = problemsById.get(ans.roomProblemId()); - if (problem == null || !problem.getRoom().getId().equals(roomId)) { - throw new BusinessException(RoomErrorCode.PROBLEM_NOT_FOUND); + if (Boolean.TRUE.equals(roomUser.getIsAttended())) { + throw new BusinessException(RoomErrorCode.ALREADY_ENDED); } - return UserRoomAnswer.of(user, problem, ans.userAnswer(), null); - }) - .toList(); - userRoomAnswerRepository.saveAll(userRoomAnswers); - // 응시 완료 상태 업데이트 - roomUser.attend(); + // 5. 동기 DB 저장 대신 Kafka 전송 후 < 20ms 즉시 응답 반환 + roomSubmitKafkaProducer.send(RoomSubmitEventMessage.of(userId, roomId, request.answers())); - log.info("[RoomService] 시험 답안 제출 완료 - userId: {}, roomId: {}, 제출 문항 수: {}", - userId, roomId, userRoomAnswers.size()); + } finally { + if (lock.isHeldByCurrentThread()) { + lock.unlock(); + } + } + } else { + log.warn("[RoomService] 답안 제출 분산 락 획득 타임아웃 - userId: {}, roomId: {}", userId, roomId); + throw new BusinessException(AuthErrorCode.LOCK_ACQUISITION_FAILED, "답안 제출 락 획득 실패"); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new BusinessException(GlobalErrorCode.INTERNAL_SERVER_ERROR); + } } @Override From 61438ebfb709fcc16b6b41e18c3f136503a9eaa9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=ED=99=8D=EC=84=B1=ED=9C=98?= Date: Tue, 11 Aug 2026 21:18:18 +0900 Subject: [PATCH 2/2] refactor(room): enhance answer submission pipeline with Kafka async processing and Redisson lock - Replace pessimistic DB lock with Redisson distributed lock outside DB transaction - Prevent TOCTOU duplicate submissions using Redis claim markers (setIfAbsent) - Enforce synchronous producer ACK wait with 3-second timeout for message durability - Implement RoomSubmitKafkaConsumer with atomic Redis SETNX deduplication and rollback cleanup - Validate problem room ownership and count synchronously in Service and Consumer - Throw NOT_ROOM_PARTICIPANT exception on missing RoomUser during attendance update - Add LOCK_ACQUISITION_FAILED error code to RoomErrorCode for domain boundary isolation - Add non-null invariant checks to RoomSubmitEventMessage compact constructor --- .../consumer/RoomSubmitKafkaConsumer.java | 67 ++++++++++++++----- .../room/event/RoomSubmitEventMessage.java | 2 + .../domain/room/exception/RoomErrorCode.java | 3 +- .../producer/RoomSubmitKafkaProducer.java | 36 ++++++---- .../domain/room/service/RoomServiceImpl.java | 35 +++++++++- 5 files changed, 109 insertions(+), 34 deletions(-) diff --git a/momogo-api/src/main/java/com/momogo/api/room/consumer/RoomSubmitKafkaConsumer.java b/momogo-api/src/main/java/com/momogo/api/room/consumer/RoomSubmitKafkaConsumer.java index 319c96a..0932949 100644 --- a/momogo-api/src/main/java/com/momogo/api/room/consumer/RoomSubmitKafkaConsumer.java +++ b/momogo-api/src/main/java/com/momogo/api/room/consumer/RoomSubmitKafkaConsumer.java @@ -1,6 +1,8 @@ package com.momogo.api.room.consumer; import com.momogo.core.common.config.KafkaTopics; +import com.momogo.core.common.exception.BusinessException; +import com.momogo.core.domain.room.exception.RoomErrorCode; import com.momogo.core.domain.room.entity.RoomProblem; import com.momogo.core.domain.room.entity.RoomUser; import com.momogo.core.domain.room.entity.RoomUserId; @@ -19,6 +21,7 @@ import java.util.stream.Collectors; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.dao.DataIntegrityViolationException; import org.springframework.data.redis.core.RedisTemplate; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; @@ -50,42 +53,70 @@ public void consume(RoomSubmitEventMessage message) { log.info("[RoomSubmitKafkaConsumer] 카프카 답안 수신 - eventId: {}, userId: {}, roomId: {}", message.eventId(), message.userId(), message.roomId()); - // 1. Redis 멱등성 키 검증 (중복 처리 방지) + // 1. Redis atomic setIfAbsent로 "PROCESSING" 상태 선점 (원자적 멱등성 검사) String dedupKey = DEDUP_KEY_PREFIX + message.eventId(); - if (Boolean.TRUE.equals(redisTemplate.hasKey(dedupKey))) { - log.warn("[RoomSubmitKafkaConsumer] 이미 처리된 답안 제출 이벤트, 중복 스킵 - eventId: {}", message.eventId()); + Boolean isAcquired = redisTemplate.opsForValue().setIfAbsent(dedupKey, "PROCESSING", Duration.ofMinutes(5)); + if (Boolean.FALSE.equals(isAcquired)) { + log.warn("[RoomSubmitKafkaConsumer] 이미 처리 중이거나 완료된 답안 제출 이벤트, 중복 스킵 - eventId: {}", message.eventId()); return; } - // 2. 추가 SELECT 없이 EntityManager 프록시 참조 생성 + // 2. DB 트랜잭션 종결 상태별 후처리 (커밋 성공 시 DONE 갱신, 예외/롤백 시 키 삭제로 카프카 재시도 보장) + TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { + @Override + public void afterCommit() { + redisTemplate.opsForValue().set(dedupKey, "DONE", DEDUP_TTL); + } + + @Override + public void afterCompletion(int status) { + if (status != STATUS_COMMITTED) { + redisTemplate.delete(dedupKey); + redisTemplate.delete("room:submit:claimed:" + message.roomId() + ":" + message.userId()); + } + } + }); + + // 3. 추가 SELECT 없이 EntityManager 프록시 참조 생성 User userProxy = entityManager.getReference(User.class, message.userId()); - // 3. 문제 목록 조회 및 답안 일괄 매핑 + // 4. 문제 목록 조회 및 해당 방 소속 검증 List problemIds = message.answers().stream() .map(ans -> ans.roomProblemId()) .toList(); Map problemsById = roomProblemRepository.findAllById(problemIds).stream() + .filter(p -> p.getRoom().getId().equals(message.roomId())) .collect(Collectors.toMap(RoomProblem::getId, Function.identity())); + if (problemsById.size() != problemIds.size()) { + log.error("[RoomSubmitKafkaConsumer] 유효하지 않거나 해당 방에 속하지 않는 문제 ID 포함, 처리 중단 - eventId: {}, roomId: {}, 요청 문제 수: {}, 유효 문제 수: {}", + message.eventId(), message.roomId(), problemIds.size(), problemsById.size()); + throw new BusinessException(RoomErrorCode.PROBLEM_NOT_FOUND); + } + List userRoomAnswers = message.answers().stream() .map(ans -> UserRoomAnswer.of(userProxy, problemsById.get(ans.roomProblemId()), ans.userAnswer(), null)) .toList(); - // 4. DB 일괄 저장 - userRoomAnswerRepository.saveAll(userRoomAnswers); - - // 5. 응시 완료 상태 업데이트 - RoomUserId roomUserId = new RoomUserId(message.roomId(), message.userId()); - roomUserRepository.findById(roomUserId).ifPresent(RoomUser::attend); + try { + // 5. DB 일괄 저장 + userRoomAnswerRepository.saveAll(userRoomAnswers); - // 6. DB 커밋 확정 후에만 Redis 멱등성 처리 완료 키 갱신 - TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { - @Override - public void afterCommit() { - redisTemplate.opsForValue().set(dedupKey, "1", DEDUP_TTL); - } - }); + // 6. 응시 완료 상태 업데이트 + RoomUserId roomUserId = new RoomUserId(message.roomId(), message.userId()); + roomUserRepository.findById(roomUserId) + .orElseThrow(() -> { + log.error("[RoomSubmitKafkaConsumer] 참여자 정보 없음, 응시 상태 갱신 불가 - userId: {}, roomId: {}", + message.userId(), message.roomId()); + return new BusinessException(RoomErrorCode.NOT_ROOM_PARTICIPANT); + }) + .attend(); + } catch (DataIntegrityViolationException e) { + log.warn("[RoomSubmitKafkaConsumer] DB 답안 중복 유니크 제약조건 충돌, 처리를 정상 스킵함 - eventId: {}, userId: {}, roomId: {}", + message.eventId(), message.userId(), message.roomId()); + return; + } log.info("[RoomSubmitKafkaConsumer] DB 저장 완료 - userId: {}, roomId: {}, 문항 수: {}", message.userId(), message.roomId(), userRoomAnswers.size()); diff --git a/momogo-core/src/main/java/com/momogo/core/domain/room/event/RoomSubmitEventMessage.java b/momogo-core/src/main/java/com/momogo/core/domain/room/event/RoomSubmitEventMessage.java index 87f60af..9ed6a7f 100644 --- a/momogo-core/src/main/java/com/momogo/core/domain/room/event/RoomSubmitEventMessage.java +++ b/momogo-core/src/main/java/com/momogo/core/domain/room/event/RoomSubmitEventMessage.java @@ -18,6 +18,8 @@ public record RoomSubmitEventMessage( Objects.requireNonNull(eventId, "eventId는 null일 수 없습니다"); Objects.requireNonNull(userId, "userId는 null일 수 없습니다"); Objects.requireNonNull(roomId, "roomId는 null일 수 없습니다"); + Objects.requireNonNull(answers, "answers는 null일 수 없습니다"); + Objects.requireNonNull(submittedAt, "submittedAt은 null일 수 없습니다"); answers = List.copyOf(answers); } diff --git a/momogo-core/src/main/java/com/momogo/core/domain/room/exception/RoomErrorCode.java b/momogo-core/src/main/java/com/momogo/core/domain/room/exception/RoomErrorCode.java index a206dea..f12aafc 100644 --- a/momogo-core/src/main/java/com/momogo/core/domain/room/exception/RoomErrorCode.java +++ b/momogo-core/src/main/java/com/momogo/core/domain/room/exception/RoomErrorCode.java @@ -19,7 +19,8 @@ public enum RoomErrorCode implements ErrorCode { REPORT_GENERATION_FAILED(4008, "REPORT_GENERATION_FAILED", HttpStatus.INTERNAL_SERVER_ERROR, "리포트 PDF 문서 생성 중 서버 내부 오류가 발생했습니다."), DUPLICATE_ANSWER_SUBMITTED(4009, "DUPLICATE_ANSWER_SUBMITTED", HttpStatus.BAD_REQUEST, "동일한 문제에 대한 중복 답안 제출은 허용되지 않습니다."), REPORT_NOT_READY(4010, "REPORT_NOT_READY", HttpStatus.BAD_REQUEST, "아직 채점이 완료되지 않아 리포트가 준비되지 않았습니다."), - AI_GRADING_IN_PROGRESS(4011, "AI_GRADING_IN_PROGRESS", HttpStatus.BAD_REQUEST, "현재 해당 시험방의 AI 채점이 진행 중입니다."); + AI_GRADING_IN_PROGRESS(4011, "AI_GRADING_IN_PROGRESS", HttpStatus.BAD_REQUEST, "현재 해당 시험방의 AI 채점이 진행 중입니다."), + LOCK_ACQUISITION_FAILED(4012, "LOCK_ACQUISITION_FAILED", HttpStatus.SERVICE_UNAVAILABLE, "동시 요청 처리를 위한 락 획득에 실패했습니다. 잠시 후 다시 시도해주세요."); private final int numeric; private final String errorKey; diff --git a/momogo-core/src/main/java/com/momogo/core/domain/room/kafka/producer/RoomSubmitKafkaProducer.java b/momogo-core/src/main/java/com/momogo/core/domain/room/kafka/producer/RoomSubmitKafkaProducer.java index 647fec9..a474733 100644 --- a/momogo-core/src/main/java/com/momogo/core/domain/room/kafka/producer/RoomSubmitKafkaProducer.java +++ b/momogo-core/src/main/java/com/momogo/core/domain/room/kafka/producer/RoomSubmitKafkaProducer.java @@ -1,15 +1,20 @@ package com.momogo.core.domain.room.kafka.producer; import com.momogo.core.common.config.KafkaTopics; +import com.momogo.core.common.exception.BusinessException; +import com.momogo.core.common.exception.GlobalErrorCode; import com.momogo.core.domain.room.event.RoomSubmitEventMessage; +import java.util.concurrent.TimeUnit; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.support.SendResult; import org.springframework.stereotype.Component; /* - * 답안 제출 이벤트를 room-submit-events 토픽에 비동기 발행하는 전담 클래스. - * 실제 DB 저장은 momogo-api의 RoomSubmitKafkaConsumer가 비동기로 처리 + * 답안 제출 이벤트를 room-submit-events 토픽에 발행하는 전담 클래스. + * 카프카 브로커의 발행 확정(ACK)을 타임아웃 내에 동기적으로 대기하여 내구성을 보장하며, + * 실제 DB 저장은 momogo-api의 RoomSubmitKafkaConsumer가 비동기로 처리한다. */ @Slf4j @Component @@ -19,15 +24,22 @@ public class RoomSubmitKafkaProducer { private final KafkaTemplate kafkaTemplate; public void send(RoomSubmitEventMessage message) { - kafkaTemplate.send(KafkaTopics.ROOM_SUBMIT_EVENTS, message.roomId().toString(), message) - .whenComplete((result, ex) -> { - if (ex != null) { - log.error("[RoomSubmitKafkaProducer] 카프카 메시지 전송 실패 - eventId: {}, roomId: {}, topic: {}", - message.eventId(), message.roomId(), KafkaTopics.ROOM_SUBMIT_EVENTS, ex); - } else { - log.info("[RoomSubmitKafkaProducer] 카프카 메시지 전송 성공 - eventId: {}, offset: {}", - message.eventId(), result.getRecordMetadata().offset()); - } - }); + try { + SendResult result = kafkaTemplate.send( + KafkaTopics.ROOM_SUBMIT_EVENTS, + message.roomId().toString(), + message + ).get(3, TimeUnit.SECONDS); + + log.info("[RoomSubmitKafkaProducer] 카프카 메시지 전송 성공 - eventId: {}, offset: {}", + message.eventId(), result.getRecordMetadata().offset()); + } catch (Exception ex) { + log.error("[RoomSubmitKafkaProducer] 카프카 메시지 전송 실패 - eventId: {}, roomId: {}, topic: {}", + message.eventId(), message.roomId(), KafkaTopics.ROOM_SUBMIT_EVENTS, ex); + if (ex instanceof InterruptedException) { + Thread.currentThread().interrupt(); + } + throw new BusinessException(GlobalErrorCode.INTERNAL_SERVER_ERROR, "답안 제출 전송 실패"); + } } } diff --git a/momogo-core/src/main/java/com/momogo/core/domain/room/service/RoomServiceImpl.java b/momogo-core/src/main/java/com/momogo/core/domain/room/service/RoomServiceImpl.java index ef089de..26a33a7 100644 --- a/momogo-core/src/main/java/com/momogo/core/domain/room/service/RoomServiceImpl.java +++ b/momogo-core/src/main/java/com/momogo/core/domain/room/service/RoomServiceImpl.java @@ -52,6 +52,7 @@ import java.io.ByteArrayOutputStream; import java.io.IOException; import java.io.InputStream; +import java.time.Duration; import java.time.OffsetDateTime; import java.util.ArrayList; import java.util.Collections; @@ -71,7 +72,9 @@ import org.redisson.api.RLock; import org.redisson.api.RedissonClient; import org.springframework.context.ApplicationEventPublisher; +import org.springframework.data.redis.core.RedisTemplate; import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Propagation; import org.springframework.transaction.annotation.Transactional; @Slf4j @@ -93,8 +96,11 @@ public class RoomServiceImpl implements RoomService{ private final AiGradingProducer aiGradingProducer; private final RedissonClient redissonClient; private final RoomSubmitKafkaProducer roomSubmitKafkaProducer; + private final RedisTemplate redisTemplate; private static final String SUBMIT_LOCK_PREFIX = "lock:room:submit:"; + private static final String SUBMIT_CLAIM_PREFIX = "room:submit:claimed:"; + private static final Duration SUBMIT_CLAIM_TTL = Duration.ofDays(7); // 검증된 타겟 객체들을 묶기 위한 private record private record ValidatedRoomTarget(Space space, List targetUsers) {} @@ -247,6 +253,7 @@ public List getRoomProblems(UUID userId, UUID roomId) { } @Override + @Transactional(propagation = Propagation.NOT_SUPPORTED) public void submitRoomAnswer(UUID userId, UUID roomId, RoomAnswerSubmitRequest request) { log.info("[RoomService] 시험 답안 제출 시도 - userId: {}, roomId: {}", userId, roomId); @@ -286,7 +293,28 @@ public void submitRoomAnswer(UUID userId, UUID roomId, RoomAnswerSubmitRequest r throw new BusinessException(RoomErrorCode.ALREADY_ENDED); } - // 5. 동기 DB 저장 대신 Kafka 전송 후 < 20ms 즉시 응답 반환 + // 4-1. 비동기 저장 완료 전 중복 제출을 즉시 차단하는 Redis Claim 마커 (Atomic SETNX) + String claimKey = SUBMIT_CLAIM_PREFIX + roomId + ":" + userId; + Boolean claimed = redisTemplate.opsForValue().setIfAbsent(claimKey, "1", SUBMIT_CLAIM_TTL); + if (!Boolean.TRUE.equals(claimed)) { + log.warn("[RoomService] 이미 답안 제출 요청이 처리 중이거나 완료된 유저 - userId: {}, roomId: {}", userId, roomId); + throw new BusinessException(RoomErrorCode.ALREADY_ENDED); + } + + // 5. 문제 존재 및 해당 시험방 소속 검증 + List problemIds = request.answers().stream() + .map(ProblemAnswerRequest::roomProblemId) + .toList(); + + long validProblemCount = roomProblemRepository.findAllById(problemIds).stream() + .filter(p -> p.getRoom().getId().equals(roomId)) + .count(); + + if (validProblemCount != problemIds.size()) { + throw new BusinessException(RoomErrorCode.PROBLEM_NOT_FOUND); + } + + // 6. 동기 DB 저장 대신 Kafka 전송 후 < 20ms 즉시 응답 반환 roomSubmitKafkaProducer.send(RoomSubmitEventMessage.of(userId, roomId, request.answers())); } finally { @@ -296,11 +324,12 @@ public void submitRoomAnswer(UUID userId, UUID roomId, RoomAnswerSubmitRequest r } } else { log.warn("[RoomService] 답안 제출 분산 락 획득 타임아웃 - userId: {}, roomId: {}", userId, roomId); - throw new BusinessException(AuthErrorCode.LOCK_ACQUISITION_FAILED, "답안 제출 락 획득 실패"); + throw new BusinessException(RoomErrorCode.LOCK_ACQUISITION_FAILED); } } catch (InterruptedException e) { + log.error("[RoomService] 답안 제출 대기 중 인터럽트 발생 - userId: {}, roomId: {}", userId, roomId, e); Thread.currentThread().interrupt(); - throw new BusinessException(GlobalErrorCode.INTERNAL_SERVER_ERROR); + throw new BusinessException(GlobalErrorCode.INTERNAL_SERVER_ERROR, "답안 제출 대기 중 인터럽트가 발생했습니다."); } }