Skip to content

Commit 19c7335

Browse files
committed
refactor(Email): 이메일 비동기 발송 요청에 Redis Queue 도입 (#1036)
refactor(Email): 이메일 비동기 발송 요청에 Redis Queue 도입 (#1036) * Create draft PR for #1010 * feat(EmailStatus): 이메일 전송 상태 Enum 추가 * feat(Email): 이메일 도메인 필드 isSucceed 제거 및 EmailStatus 추가 * feat: Email 테이블 스키마 변경 적용을 위한 Flyway 스크립트 작성 * feat(Email): Redis 기반 이메일 큐 서비스 구현 * feat(Email): 이메일 중복 요청 예외 처리 구현 * feat(Email): 이메일 발송 시스템에 Redis 큐 적용 * refactor: 멤버 테스트 코드 스타일 개선 * feat(EmailQueueService): Redis Queue에 메일 발송 처리 중 상태 추가 및 테스트 생성 * fix: 불필요 수정사항 롤백 * refactor: Email Status(발송 결과 상태값) API 응답 프론트 작업에 맞춰 변경 롤백 * fix-be: 스케줄러를 이용한 워커 활성화 및 DTO 수정 * fix-be: flyway v2.4 version 충돌 파일명 수정
1 parent c4cf2b2 commit 19c7335

18 files changed

Lines changed: 549 additions & 72 deletions

File tree

backend/src/main/java/com/cruru/CruruApplication.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,8 +3,10 @@
33
import org.springframework.boot.SpringApplication;
44
import org.springframework.boot.autoconfigure.SpringBootApplication;
55
import org.springframework.boot.context.properties.ConfigurationPropertiesScan;
6+
import org.springframework.scheduling.annotation.EnableScheduling;
67

78
@SpringBootApplication
9+
@EnableScheduling
810
@ConfigurationPropertiesScan
911
public class CruruApplication {
1012

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,11 @@
1+
package com.cruru.advice.badrequest;
2+
3+
import com.cruru.advice.CruruCustomException;
4+
import org.springframework.http.HttpStatus;
5+
6+
public class TooManyRequestException extends CruruCustomException {
7+
8+
public TooManyRequestException(String message) {
9+
super(message, HttpStatus.TOO_MANY_REQUESTS);
10+
}
11+
}

backend/src/main/java/com/cruru/email/domain/Email.java

Lines changed: 12 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,8 @@
77
import com.cruru.email.exception.EmailSubjectLengthException;
88
import jakarta.persistence.Column;
99
import jakarta.persistence.Entity;
10+
import jakarta.persistence.EnumType;
11+
import jakarta.persistence.Enumerated;
1012
import jakarta.persistence.FetchType;
1113
import jakarta.persistence.GeneratedValue;
1214
import jakarta.persistence.GenerationType;
@@ -46,17 +48,18 @@ public class Email extends BaseEntity {
4648
@Column(columnDefinition = "TEXT")
4749
private String content;
4850

49-
@Column(name = "is_succeed")
50-
private Boolean isSucceed;
51+
@Column(name = "email_status")
52+
@Enumerated(EnumType.STRING)
53+
private EmailStatus status;
5154

52-
public Email(Club from, Applicant to, String subject, String content, Boolean isSucceed) {
55+
public Email(Club from, Applicant to, String subject, String content, EmailStatus status) {
5356
validateSubjectLength(subject);
5457
validateContentLength(content);
5558
this.from = from;
5659
this.to = to;
5760
this.subject = subject;
5861
this.content = content;
59-
this.isSucceed = isSucceed;
62+
this.status = status;
6063
}
6164

6265
private void validateSubjectLength(String subject) {
@@ -71,6 +74,10 @@ private void validateContentLength(String content) {
7174
}
7275
}
7376

77+
public void updateStatus(EmailStatus status) {
78+
this.status = status;
79+
}
80+
7481
@Override
7582
public boolean equals(Object o) {
7683
if (this == o) {
@@ -96,7 +103,7 @@ public String toString() {
96103
", to=" + to +
97104
", subject='" + subject + '\'' +
98105
", content='" + content + '\'' +
99-
", isSucceed=" + isSucceed +
106+
", isSucceed=" + status +
100107
'}';
101108
}
102109
}
Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,8 @@
1+
package com.cruru.email.domain;
2+
3+
public enum EmailStatus {
4+
PENDING,
5+
DELIVERED,
6+
FAILED
7+
;
8+
}
Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,22 @@
1+
package com.cruru.email.dto;
2+
3+
import java.io.Serializable;
4+
import java.util.List;
5+
import lombok.AllArgsConstructor;
6+
import lombok.Data;
7+
import lombok.Getter;
8+
import lombok.NoArgsConstructor;
9+
10+
@Data
11+
@NoArgsConstructor
12+
@AllArgsConstructor
13+
@Getter
14+
public class EmailQueueMessage implements Serializable {
15+
16+
private Long clubId;
17+
private Long applicantId;
18+
private String toEmail;
19+
private String subject;
20+
private String content;
21+
private List<String> attachmentPaths;
22+
}
Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,11 @@
1+
package com.cruru.email.exception;
2+
3+
import com.cruru.advice.CruruCustomException;
4+
import org.springframework.http.HttpStatus;
5+
6+
public class EmailRedisException extends CruruCustomException {
7+
8+
public EmailRedisException(String message) {
9+
super(message, HttpStatus.INTERNAL_SERVER_ERROR);
10+
}
11+
}
Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,12 @@
1+
package com.cruru.email.exception.badrequest;
2+
3+
import com.cruru.advice.badrequest.TooManyRequestException;
4+
5+
public class EmailDuplicatedRequestException extends TooManyRequestException {
6+
7+
private static final String MESSAGE = ": 발송 대기 중인 이메일 요청과 중복됩니다.";
8+
9+
public EmailDuplicatedRequestException(String emailAddress) {
10+
super(emailAddress + MESSAGE);
11+
}
12+
}

backend/src/main/java/com/cruru/email/facade/EmailFacade.java

Lines changed: 24 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -10,9 +10,13 @@
1010
import com.cruru.email.controller.response.EmailHistoryResponse;
1111
import com.cruru.email.controller.response.EmailHistoryResponses;
1212
import com.cruru.email.domain.Email;
13+
import com.cruru.email.domain.EmailStatus;
14+
import com.cruru.email.dto.EmailQueueMessage;
1315
import com.cruru.email.exception.EmailAttachmentsException;
1416
import com.cruru.email.exception.EmailConflictException;
17+
import com.cruru.email.exception.badrequest.EmailDuplicatedRequestException;
1518
import com.cruru.email.service.EmailKeywordConverter;
19+
import com.cruru.email.service.EmailQueueService;
1620
import com.cruru.email.service.EmailRedisClient;
1721
import com.cruru.email.service.EmailService;
1822
import com.cruru.email.util.FileUtil;
@@ -36,26 +40,34 @@ public class EmailFacade {
3640
private final MemberService memberService;
3741
private final EmailRedisClient emailRedisClient;
3842
private final EmailKeywordConverter emailKeywordConverter;
43+
private final EmailQueueService emailQueueService;
3944

4045
public void send(EmailRequest request) {
4146
Club from = clubService.findById(request.clubId());
4247
List<Applicant> applicants = applicantService.findAllByIds(request.applicantIds());
43-
sendAndSave(from, applicants, request.subject(), request.content(), request.files());
48+
sendAndQueue(from, applicants, request.subject(), request.content(), request.files());
4449
}
4550

46-
private void sendAndSave(Club from, List<Applicant> tos, String subject, String text, List<MultipartFile> files) {
51+
private void sendAndQueue(Club from, List<Applicant> tos, String subject, String text, List<MultipartFile> files) {
4752
List<File> tempFiles = saveTempFiles(from, subject, files);
53+
List<String> attachmentPaths = tempFiles.stream().map(File::getAbsolutePath).toList();
4854

49-
List<CompletableFuture<Void>> futures = tos.stream()
50-
.map(to -> {
51-
String content = emailKeywordConverter.convert(text, from, to);
52-
return emailService.send(from, to, subject, content, tempFiles);
53-
})
54-
.map(future -> future.thenAccept(emailService::save))
55-
.toList();
5655

57-
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0]))
58-
.thenRun(() -> FileUtil.deleteFiles(tempFiles));
56+
tos.forEach(to -> {
57+
// 이메일 내용에 대해 동적 키워드 변환을 미리 적용
58+
String convertedContent = emailKeywordConverter.convert(text, from, to);
59+
// EmailQueueMessage 생성
60+
EmailQueueMessage message = new EmailQueueMessage(
61+
from.getId(),
62+
to.getId(),
63+
to.getEmail(),
64+
subject,
65+
convertedContent,
66+
attachmentPaths
67+
);
68+
// 큐에 등록 (이미 등록된 요청은 중복으로 등록되지 않음)
69+
emailQueueService.enqueueEmail(message);
70+
});
5971
}
6072

6173
private List<File> saveTempFiles(Club from, String subject, List<MultipartFile> files) {
@@ -103,7 +115,7 @@ private EmailHistoryResponse toEmailResponse(Email email) {
103115
email.getSubject(),
104116
email.getContent(),
105117
email.getCreatedDate(),
106-
email.getIsSucceed()
118+
email.getStatus().equals(EmailStatus.DELIVERED)
107119
);
108120
}
109121
}
Lines changed: 136 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,136 @@
1+
package com.cruru.email.service;
2+
3+
import com.cruru.applicant.domain.Applicant;
4+
import com.cruru.applicant.service.ApplicantService;
5+
import com.cruru.club.domain.Club;
6+
import com.cruru.club.service.ClubService;
7+
import com.cruru.email.dto.EmailQueueMessage;
8+
import com.cruru.email.exception.EmailRedisException;
9+
import com.fasterxml.jackson.databind.ObjectMapper;
10+
import java.io.File;
11+
import java.nio.charset.StandardCharsets;
12+
import java.security.MessageDigest;
13+
import java.util.List;
14+
import java.util.concurrent.TimeUnit;
15+
import lombok.RequiredArgsConstructor;
16+
import lombok.extern.slf4j.Slf4j;
17+
import org.springframework.data.redis.core.RedisTemplate;
18+
import org.springframework.scheduling.annotation.Scheduled;
19+
import org.springframework.stereotype.Service;
20+
21+
@Service
22+
@RequiredArgsConstructor
23+
@Slf4j
24+
public class EmailQueueService {
25+
26+
private static final String EMAIL_REQUEST_KEY_PREFIX = "email_send:";
27+
private static final String EMAIL_QUEUE_KEY = "email_queue";
28+
private static final String EMAIL_PROCESSING_KEY = "email_queue_processing";
29+
private static final long IDEMPOTENCY_TTL_MINUTES = 60;
30+
31+
private final RedisTemplate<String, String> redisTemplate;
32+
private final ObjectMapper objectMapper;
33+
private final EmailService emailService;
34+
private final ClubService clubService;
35+
private final ApplicantService applicantService;
36+
37+
public boolean enqueueEmail(EmailQueueMessage message) {
38+
try {
39+
String uniqueKey = generateUniqueKey(message.getToEmail(), message.getSubject(), message.getContent());
40+
String redisUniqueKey = EMAIL_REQUEST_KEY_PREFIX + uniqueKey;
41+
42+
// 중복 여부 체크: setIfAbsent가 true이면 최초 등록, false면 중복
43+
Boolean success = redisTemplate.opsForValue().setIfAbsent(
44+
redisUniqueKey,
45+
"1",
46+
IDEMPOTENCY_TTL_MINUTES,
47+
TimeUnit.MINUTES
48+
);
49+
if (success == null || !success) {
50+
log.info("중복 이메일 발송 요청 감지: {}", redisUniqueKey);
51+
return false;
52+
}
53+
// JSON으로 변환 후, Redis 리스트 큐에 적재
54+
String messageJson = objectMapper.writeValueAsString(message);
55+
redisTemplate.opsForList().rightPush(EMAIL_QUEUE_KEY, messageJson);
56+
log.info("이메일 요청 큐에 등록: {}", messageJson);
57+
return true;
58+
} catch (Exception e) {
59+
log.error("이메일 요청 등록 실패", e);
60+
return false;
61+
}
62+
}
63+
64+
private String generateUniqueKey(String toEmail, String subject, String content) {
65+
try {
66+
MessageDigest md = MessageDigest.getInstance("MD5");
67+
String combined = toEmail + subject + content;
68+
byte[] digest = md.digest(combined.getBytes(StandardCharsets.UTF_8));
69+
StringBuilder hexString = new StringBuilder();
70+
for (byte b : digest) {
71+
String hex = Integer.toHexString(0xff & b);
72+
if (hex.length() == 1) {
73+
hexString.append('0');
74+
}
75+
hexString.append(hex);
76+
}
77+
return hexString.toString();
78+
} catch (Exception e) {
79+
throw new EmailRedisException("Redis 고유 키 생성 중 오류");
80+
}
81+
}
82+
83+
@Scheduled(fixedDelay = 5000)
84+
public void processQueue() {
85+
try {
86+
// 매 스케줄마다 처리 중인 메일 복구 시도
87+
recoverStuckMessages();
88+
89+
String messageJson;
90+
// leftPop 대신 rightPopAndLeftPush를 사용하여 원자적으로 처리 중 상태로 이동
91+
while ((messageJson = redisTemplate.opsForList().rightPopAndLeftPush(EMAIL_QUEUE_KEY, EMAIL_PROCESSING_KEY))
92+
!= null) {
93+
log.info("큐에서 이메일 요청 처리 시작: {}", messageJson);
94+
EmailQueueMessage message = objectMapper.readValue(messageJson, EmailQueueMessage.class);
95+
Club from = clubService.findById(message.getClubId());
96+
Applicant to = applicantService.findById(message.getApplicantId());
97+
List<File> attachments = null;
98+
if (message.getAttachmentPaths() != null) {
99+
attachments = message.getAttachmentPaths()
100+
.stream()
101+
.map(File::new)
102+
.toList();
103+
}
104+
105+
// 이메일 전송: send()는 비동기로 처리 후 CompletableFuture<Email> 반환
106+
String finalMessageJson = messageJson;
107+
emailService.send(from, to, message.getSubject(), message.getContent(), attachments)
108+
.thenAccept(email -> {
109+
// 성공적으로 처리되면 처리 중 리스트에서 제거
110+
redisTemplate.opsForList().remove(EMAIL_PROCESSING_KEY, 1, finalMessageJson);
111+
emailService.save(email);
112+
log.info("이메일 전송 완료 및 처리 대기열에서 제거");
113+
}).exceptionally(ex -> {
114+
log.error("이메일 전송 중 오류 발생: {}", ex.getMessage());
115+
// 오류 발생 시 처리 중 리스트에서 제거 후 다시 메인 큐로 이동
116+
redisTemplate.opsForList().remove(EMAIL_PROCESSING_KEY, 1, finalMessageJson);
117+
redisTemplate.opsForList().rightPush(EMAIL_QUEUE_KEY, finalMessageJson);
118+
log.info("전송 실패한 메일을 재처리 큐에 등록");
119+
return null;
120+
});
121+
}
122+
} catch (Exception e) {
123+
log.error("이메일 Redis 큐 처리 중 예외 발생", e);
124+
}
125+
}
126+
127+
private void recoverStuckMessages() {
128+
int count = 0;
129+
while (redisTemplate.opsForList().rightPopAndLeftPush(EMAIL_PROCESSING_KEY, EMAIL_QUEUE_KEY) != null) {
130+
count++;
131+
}
132+
if (count > 0) {
133+
log.info("애플리케이션 시작 시 처리되지 않은 메일 {}건 복원 완료", count);
134+
}
135+
}
136+
}

backend/src/main/java/com/cruru/email/service/EmailService.java

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
import com.cruru.applicant.domain.Applicant;
44
import com.cruru.club.domain.Club;
55
import com.cruru.email.domain.Email;
6+
import com.cruru.email.domain.EmailStatus;
67
import com.cruru.email.domain.repository.EmailRepository;
78
import com.cruru.email.util.EmailTemplate;
89
import jakarta.mail.MessagingException;
@@ -31,6 +32,10 @@ public class EmailService {
3132
@Async
3233
public CompletableFuture<Email> send(
3334
Club from, Applicant to, String subject, String content, List<File> tempFiles) {
35+
36+
Email email = new Email(from, to, subject, content, EmailStatus.PENDING);
37+
emailRepository.save(email);
38+
3439
try {
3540
MimeMessage message = mailSender.createMimeMessage();
3641
MimeMessageHelper helper = new MimeMessageHelper(message, true, "UTF-8");
@@ -42,11 +47,13 @@ public CompletableFuture<Email> send(
4247
}
4348
mailSender.send(message);
4449

50+
email.updateStatus(EmailStatus.DELIVERED);
4551
log.info("이메일 전송 성공: from={}, to={}, subject={}", from.getId(), to.getEmail(), subject);
46-
return CompletableFuture.completedFuture(new Email(from, to, subject, content, true));
52+
return CompletableFuture.completedFuture(email);
4753
} catch (MessagingException | MailException e) {
48-
log.info("이메일 전송 실패: from={}, to={}, subject={}", from.getId(), to.getEmail(), e.getMessage());
49-
return CompletableFuture.completedFuture(new Email(from, to, subject, content, false));
54+
log.info("이메일 전송 실패: from={}, to={}, subject={}, error={}", from.getId(), to.getEmail(), subject, e.getMessage());
55+
email.updateStatus(EmailStatus.FAILED);
56+
return CompletableFuture.completedFuture(email);
5057
}
5158
}
5259

0 commit comments

Comments
 (0)