Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,123 @@
package com.momogo.api.room.consumer;

import com.momogo.core.common.config.KafkaTopics;
import com.momogo.core.common.config.RoomRedisKeys;
import com.momogo.core.common.exception.BusinessException;
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.exception.RoomErrorCode;
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 Duration DEDUP_TTL = Duration.ofHours(24);

private final UserRoomAnswerRepository userRoomAnswerRepository;
private final RoomUserRepository roomUserRepository;
private final RoomProblemRepository roomProblemRepository;
private final RedisTemplate<String, Object> 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 = RoomRedisKeys.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(RoomRedisKeys.SUBMIT_CLAIM_PREFIX + message.roomId() + ":" + message.userId());
}
}
});

// 3. 추가 SELECT 없이 EntityManager 프록시 참조 생성
User userProxy = entityManager.getReference(User.class, message.userId());

// 4. 문제 목록 조회 및 해당 방 소속 검증
List<UUID> problemIds = message.answers().stream()
.map(ans -> ans.roomProblemId())
.toList();

Map<UUID, RoomProblem> 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<UserRoomAnswer> userRoomAnswers = message.answers().stream()
.map(ans -> UserRoomAnswer.of(userProxy, problemsById.get(ans.roomProblemId()), ans.userAnswer(), null))
.toList();

// 5. 참여자 자격 조회 및 중복 응시 사전 검증 (트랜잭션 rollback-only 예외 방지)
RoomUserId roomUserId = new RoomUserId(message.roomId(), message.userId());
RoomUser roomUser = roomUserRepository.findById(roomUserId)
.orElseThrow(() -> {
log.error("[RoomSubmitKafkaConsumer] 참여자 정보 없음, 처리 중단 - userId: {}, roomId: {}",
message.userId(), message.roomId());
return new BusinessException(RoomErrorCode.NOT_ROOM_PARTICIPANT);
});

if (Boolean.TRUE.equals(roomUser.getIsAttended())) {
log.warn("[RoomSubmitKafkaConsumer] 이미 응시 완료된 유저의 중복 제출 이벤트, 처리 스킵 - eventId: {}, userId: {}, roomId: {}",
message.eventId(), message.userId(), message.roomId());
return;
}

// 6. DB 일괄 저장 및 응시 처리
userRoomAnswerRepository.saveAll(userRoomAnswers);
roomUser.attend();

log.info("[RoomSubmitKafkaConsumer] DB 저장 완료 - userId: {}, roomId: {}, 문항 수: {}",
message.userId(), message.roomId(), userRoomAnswers.size());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
package com.momogo.core.common.config;

/**
* Room 도메인 관련 Redis 키 상수.
* 서비스와 컨슈머 간 동일한 Redis 키를 안전하게 공유하기 위해 관리한다.
*/
public final class RoomRedisKeys {

public static final String SUBMIT_LOCK_PREFIX = "lock:room:submit:";
public static final String SUBMIT_CLAIM_PREFIX = "room:submit:claimed:";
public static final String DEDUP_KEY_PREFIX = "room:submit:processed:";

private RoomRedisKeys() {
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
package com.momogo.core.common.lock;

import com.momogo.core.common.exception.BusinessException;
import com.momogo.core.common.exception.GlobalErrorCode;
import com.momogo.core.domain.room.exception.RoomErrorCode;
import java.util.concurrent.TimeUnit;
import java.util.function.Supplier;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.redisson.api.RLock;
import org.redisson.api.RedissonClient;
import org.springframework.stereotype.Component;

/**
* Redisson 분산 락 획득 및 해제를 담당하는 헬퍼 컴포넌트.
* 분산 락 처리 중 발생하는 락 타임아웃, 인터럽트 및 언락 처리를 공통화한다.
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class DistributedLockExecutor {

private final RedissonClient redissonClient;

public void executeWithLock(String lockKey, long waitSeconds, Runnable action) {
executeWithLock(lockKey, waitSeconds, () -> {
action.run();
return null;
});
}

public <T> T executeWithLock(String lockKey, long waitSeconds, Supplier<T> supplier) {
RLock lock = redissonClient.getLock(lockKey);
try {
if (!lock.tryLock(waitSeconds, TimeUnit.SECONDS)) {
log.warn("[DistributedLockExecutor] 분산 락 획득 실패 - lockKey: {}", lockKey);
throw new BusinessException(RoomErrorCode.LOCK_ACQUISITION_FAILED);
}
try {
return supplier.get();
} finally {
if (lock.isHeldByCurrentThread()) {
lock.unlock();
}
}
} catch (InterruptedException e) {
log.error("[DistributedLockExecutor] 분산 락 대기 중 인터럽트 발생 - lockKey: {}", lockKey, e);
Thread.currentThread().interrupt();
throw new BusinessException(GlobalErrorCode.INTERNAL_SERVER_ERROR, "락 대기 중 인터럽트가 발생했습니다.");
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
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;

import java.time.ZoneOffset;

public record RoomSubmitEventMessage(
UUID eventId, // 카프카 재전달 시 중복 처리 방지용 식별자 (Redis Dedup 키)
UUID userId,
UUID roomId,
List<ProblemAnswerRequest> 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<ProblemAnswerRequest> answers) {
return new RoomSubmitEventMessage(UUID.randomUUID(), userId, roomId, answers, OffsetDateTime.now(ZoneOffset.UTC));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,9 @@ 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, "동시 요청 처리를 위한 락 획득에 실패했습니다. 잠시 후 다시 시도해주세요."),
ALREADY_SUBMITTED(4013, "ALREADY_SUBMITTED", HttpStatus.CONFLICT, "이미 답안을 제출한 시험입니다.");

private final int numeric;
private final String errorKey;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
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.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
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<String, Object> kafkaTemplate;

public void send(RoomSubmitEventMessage message) {
String partitionKey = message.roomId() + ":" + message.userId();
try {
SendResult<String, Object> result = kafkaTemplate.send(
KafkaTopics.ROOM_SUBMIT_EVENTS,
partitionKey,
message
).get(3, TimeUnit.SECONDS);

log.info("[RoomSubmitKafkaProducer] 카프카 메시지 전송 성공 - eventId: {}, offset: {}",
message.eventId(), result.getRecordMetadata().offset());
} catch (InterruptedException ex) {
Thread.currentThread().interrupt();
log.error("[RoomSubmitKafkaProducer] 카프카 메시지 전송 대기 중 인터럽트 - eventId: {}, roomId: {}",
message.eventId(), message.roomId(), ex);
throw new BusinessException(GlobalErrorCode.INTERNAL_SERVER_ERROR, "답안 제출 전송 실패");
} catch (ExecutionException | TimeoutException ex) {
log.error("[RoomSubmitKafkaProducer] 카프카 메시지 전송 실패 - eventId: {}, roomId: {}, topic: {}",
message.eventId(), message.roomId(), KafkaTopics.ROOM_SUBMIT_EVENTS, ex);
throw new BusinessException(GlobalErrorCode.INTERNAL_SERVER_ERROR, "답안 제출 전송 실패");
}
}
}
Loading