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..0932949 --- /dev/null +++ b/momogo-api/src/main/java/com/momogo/api/room/consumer/RoomSubmitKafkaConsumer.java @@ -0,0 +1,124 @@ +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; +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.dao.DataIntegrityViolationException; +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 atomic setIfAbsent로 "PROCESSING" 상태 선점 (원자적 멱등성 검사) + String dedupKey = DEDUP_KEY_PREFIX + 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. 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()); + + // 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(); + + try { + // 5. DB 일괄 저장 + userRoomAnswerRepository.saveAll(userRoomAnswers); + + // 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/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..9ed6a7f --- /dev/null +++ b/momogo-core/src/main/java/com/momogo/core/domain/room/event/RoomSubmitEventMessage.java @@ -0,0 +1,29 @@ +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일 수 없습니다"); + Objects.requireNonNull(answers, "answers는 null일 수 없습니다"); + Objects.requireNonNull(submittedAt, "submittedAt은 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/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 new file mode 100644 index 0000000..a474733 --- /dev/null +++ b/momogo-core/src/main/java/com/momogo/core/domain/room/kafka/producer/RoomSubmitKafkaProducer.java @@ -0,0 +1,45 @@ +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 토픽에 발행하는 전담 클래스. + * 카프카 브로커의 발행 확정(ACK)을 타임아웃 내에 동기적으로 대기하여 내구성을 보장하며, + * 실제 DB 저장은 momogo-api의 RoomSubmitKafkaConsumer가 비동기로 처리한다. + */ +@Slf4j +@Component +@RequiredArgsConstructor +public class RoomSubmitKafkaProducer { + + private final KafkaTemplate kafkaTemplate; + + public void send(RoomSubmitEventMessage message) { + 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 36ad8c7..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; @@ -63,8 +64,17 @@ 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.data.redis.core.RedisTemplate; import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Propagation; import org.springframework.transaction.annotation.Transactional; @Slf4j @@ -84,6 +94,13 @@ 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 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) {} @@ -236,11 +253,11 @@ public List getRoomProblems(UUID userId, UUID roomId) { } @Override - @Transactional + @Transactional(propagation = Propagation.NOT_SUPPORTED) 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 +265,72 @@ 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)); + // 2. Redisson 분산 락 획득 (Redisson Watchdog 활용을 위해 leaseTime 생략) + String lockKey = SUBMIT_LOCK_PREFIX + roomId + ":" + userId; + RLock lock = redissonClient.getLock(lockKey); - // 시간 검증 - 시험 시작 전 제출 시도 차단 - OffsetDateTime now = OffsetDateTime.now(); - if (now.isBefore(room.getTestStartAt())) { - throw new BusinessException(RoomErrorCode.INVALID_ACCESS_BEFORE_START); - } + try { + if (lock.tryLock(3, TimeUnit.SECONDS)) { + try { + // 3. 비관적 락(findByIdForUpdate) 대신 락이 없는 일반 조회 사용 + Room room = findRoomOrThrow(roomId); - // 상태 검증 - 이미 최종적으로 끝난 시험이면 추가 제출 불가 - if (Boolean.TRUE.equals(room.getIsEnded())) { - throw new BusinessException(RoomErrorCode.ALREADY_ENDED); - } + // 4. 시간 및 응시 자격 빠른 검증 + OffsetDateTime now = OffsetDateTime.now(); + if (now.isBefore(room.getTestStartAt())) { + throw new BusinessException(RoomErrorCode.INVALID_ACCESS_BEFORE_START); + } - // 응시 대상 유저 자격 검증 - RoomUserId roomUserId = new RoomUserId(roomId, userId); - RoomUser roomUser = roomUserRepository.findByIdForUpdate(roomUserId) - .orElseThrow(() -> new BusinessException(RoomErrorCode.NOT_ROOM_PARTICIPANT)); + if (Boolean.TRUE.equals(room.getIsEnded())) { + throw new BusinessException(RoomErrorCode.ALREADY_ENDED); + } - // 요청 유저 획득 - User user = userRepository.findById(userId) - .orElseThrow(() -> new BusinessException(SpaceErrorCode.SPACE_USER_NOT_FOUND)); + RoomUserId roomUserId = new RoomUserId(roomId, userId); + RoomUser roomUser = roomUserRepository.findById(roomUserId) + .orElseThrow(() -> new BusinessException(RoomErrorCode.NOT_ROOM_PARTICIPANT)); - 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(roomUser.getIsAttended())) { + throw new BusinessException(RoomErrorCode.ALREADY_ENDED); + } - if (Boolean.TRUE.equals(roomUser.getIsAttended())) { - throw new BusinessException(RoomErrorCode.ALREADY_ENDED); - } + // 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(); - // 각 문제 답안 일괄 매핑, 저장 - List userRoomAnswers = request.answers().stream() - .map(ans -> { - RoomProblem problem = problemsById.get(ans.roomProblemId()); - if (problem == null || !problem.getRoom().getId().equals(roomId)) { + 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); } - return UserRoomAnswer.of(user, problem, ans.userAnswer(), null); - }) - .toList(); - userRoomAnswerRepository.saveAll(userRoomAnswers); - // 응시 완료 상태 업데이트 - roomUser.attend(); + // 6. 동기 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(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, "답안 제출 대기 중 인터럽트가 발생했습니다."); + } } @Override