-
Notifications
You must be signed in to change notification settings - Fork 0
RabbitMQ
이 페이지는 Notion(AIOT-Team3-InsightON)에서 정리해 옮긴 내용입니다. 최신 원본은 Notion에서 확인하세요. (일회성 정리 — 자동 동기화 아님)
⚠️ 원본 페이지에 실제 공용 서버로 보이는host/password값이 예시 코드에 그대로 적혀 있어, 이 문서에서는<REDACTED>로 가려뒀습니다. 실제 값이 필요하면 Notion 원본 또는 팀 노션 담당자에게 확인하세요.
예시 코드는 원본 작성자가 맡았던 "쿠폰" 도메인 기준입니다. InsightOn 프로젝트에서는 telemetry.{groups_id} 라우팅으로 대체됩니다 (Rule Engine 참고 문서의 "3. RabbitMQ 라우팅 및 필터링 설계" 참고).
**메시지 브로커(Message Broker)**라고 부르는 미들웨어 소프트웨어. 시스템과 시스템 사이에서 데이터(메시지)를 대신 전달해주는 역할을 합니다.
- 없을 때: 유저가 회원가입 버튼을 누름 → 서버가 DB 저장 → 쿠폰 발급 로직(3초) → 이메일 발송(2초) → 총 5초 뒤에 "가입 완료" (유저 답답함)
- RabbitMQ 사용: 유저 가입 → RabbitMQ에 "쿠폰 줘" 메시지 던짐 → 바로 "가입 완료" 표시 → 쿠폰은 뒷단에서 알아서 처리
- 없을 때: 쿠폰 서버가 죽으면 회원가입도 같이 에러
- RabbitMQ 사용: 쿠폰 서버가 죽어도 큐에 메시지가 안전하게 쌓임. 회원가입은 정상 진행. 나중에 쿠폰 서버를 복구하면 쌓인 메시지를 처리
갑자기 대량 요청이 몰려도 RabbitMQ가 댐 역할을 해서, 쿠폰 서버는 자기가 처리할 수 있는 속도로 하나씩 가져와 처리하면 되니 서버가 터지지 않습니다.
원본에 있던 내부 구조 다이어그램(Producer/Exchange/Queue/Binding, Direct/Topic Exchange 그림)은 임시 서명 URL이라 여기 옮기지 못했습니다. Notion 원본에서 확인해주세요.
메시지를 생성하고 발송하는 애플리케이션. 절대로 Queue에 직접 넣지 않고, Exchange에게 넘기기만 합니다.
Producer가 보낸 메시지를 받아서 규칙에 따라 Queue에 분배하는 라우터. 규칙은 라우팅 키(Routing Key)를 보고 결정합니다.
메시지가 실제로 저장되는 버퍼(메모리/디스크). Consumer가 가져갈 때까지 쌓아둡니다.
Exchange와 Queue를 연결해주는 연결 고리.
-
Direct Exchange: 라우팅 키가 정확히 일치하는 큐로만 보냄 (예:
coupon.welcome→ 딱 그 큐로만). 1:1 전송, 명확한 타겟팅. -
Topic Exchange (패턴 매칭): 라우팅 키에
*,#패턴 사용 가능 (예:log.*→log.error,log.info모두 수신). 유연한 라우팅에 사용. - Fanout Exchange: 라우팅 키를 무시하고 연결된 모든 큐에 동일 메시지 복사 (전체 공지 등).
- Headers Exchange: 라우팅 키 대신 헤더 정보로 결정.
Direct / Topic이 자주 쓰이며, InsightOn에서는 Topic Exchange(telemetry.{groups_id}, telemetry.# 와일드카드 바인딩)를 사용합니다.
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>rabbitmq:
host: <REDACTED>
port: 5672
username: <REDACTED>
password: <REDACTED>user측에서 보내는 메시지는 모두 book2.dev.user.exchange를 통해 들어가고, exchange가 RoutingKey를 확인해 해당 Queue로 라우팅합니다.
@Configuration
public class RabbitConfig {
// User가 보낸 메시지를 처리할 exchange
public static final String EXCHANGE = "book2.dev.user.exchange";
// 라우팅 키가 coupon.welcome이면 이 키와 연결된 큐로 메시지 저장
public static final String ROUTING_KEY_WELCOME = "coupon.welcome";
// 라우팅 키가 coupon.birthday면 이 키와 연결된 큐로 메시지 저장
public static final String ROUTING_KEY_BIRTHDAY = "coupon.birthday";
// Provider는 Queue로 직접 접근하지 않고 Exchange 설정만 합니다.
// DirectExchange: RoutingKey가 정확히 일치해야 Queue로 데이터 적재
@Bean
public DirectExchange exchange() {
return new DirectExchange(EXCHANGE);
}
}회원가입 시 웰컴쿠폰 지급을 위해, 가입과 동시에 userId를 메시지로 보냅니다.
@Service
@RequiredArgsConstructor
public class UserService {
private final RabbitTemplate rabbitTemplate;
private final UserRepository userRepository;
public UserResponseDto registerUser(UserSignupRequest dto) {
User savedUser = userRepository.save(dto.toEntity());
// convertAndSend(exchange, routingKey, message)
rabbitTemplate.convertAndSend(
RabbitConfig.EXCHANGE,
RabbitConfig.ROUTING_KEY_WELCOME,
savedUser.getId()
);
return UserResponseDto.from(savedUser);
}
}Consumer가 메시지 처리 중 예외가 발생하면 메시지는 Queue 맨 앞으로 재삽입되어 계속 실패하며 Block 상태에 걸릴 수 있습니다. 이를 막기 위해 DLX/DLQ를 사용합니다.
실패한 메시지가 도착하는 Exchange. 실패한 메시지를 받아서 연결된 DLQ로 보내는 역할.
실패한 메시지가 최종적으로 쌓이는 Queue. 디버깅(왜 에러가 났는지 확인)과 재처리(Replay)에 사용.
@Configuration
public class RabbitConfig {
public static final String EXCHANGE = "book2.dev.user.exchange";
public static final String QUEUE_WELCOME = "book2.dev.welcome.queue";
public static final String ROUTING_KEY_WELCOME = "coupon.welcome";
public static final String QUEUE_BIRTHDAY = "book2.dev.birthday.queue";
public static final String ROUTING_KEY_BIRTHDAY = "coupon.birthday";
public static final String DLX_EXCHANGE = "book2.dev.dlx.coupon.exchange";
public static final String QUEUE_WELCOME_DLQ = "book2.dev.welcome.dlq";
public static final String DLX_ROUTING_KEY_WELCOME = "coupon.welcome.dlq";
public static final String QUEUE_BIRTHDAY_DLQ = "book2.dev.birthday.dlq";
public static final String DLX_ROUTING_KEY_BIRTHDAY = "coupon.birthday.dlq";
@Bean
public DirectExchange exchange() {
return new DirectExchange(EXCHANGE);
}
// durable(true): 서버가 다운돼도 디스크에 남아 사라지지 않음
@Bean
public Queue welcomeQueue() {
return QueueBuilder.durable(QUEUE_WELCOME)
.deadLetterExchange(DLX_EXCHANGE)
.deadLetterRoutingKey(DLX_ROUTING_KEY_WELCOME)
.build();
}
@Bean
public Binding welcomeBinding() {
return BindingBuilder.bind(welcomeQueue())
.to(exchange())
.with(ROUTING_KEY_WELCOME);
}
@Bean
public Queue welcomeDlq() {
return new Queue(QUEUE_WELCOME_DLQ, true);
}
@Bean
public Binding welcomeDlqBinding() {
return BindingBuilder.bind(welcomeDlq())
.to(dlxExchange())
.with(DLX_ROUTING_KEY_WELCOME);
}
@Bean
public Queue birthdayQueue() {
return QueueBuilder.durable(QUEUE_BIRTHDAY)
.deadLetterExchange(DLX_EXCHANGE)
.deadLetterRoutingKey(DLX_ROUTING_KEY_BIRTHDAY)
.build();
}
@Bean
public Binding birthdayBinding() {
return BindingBuilder.bind(birthdayQueue())
.to(exchange())
.with(ROUTING_KEY_BIRTHDAY);
}
@Bean
public Queue birthdayDlq() {
return new Queue(QUEUE_BIRTHDAY_DLQ, true);
}
@Bean
public Binding birthdayDlqBinding() {
return BindingBuilder.bind(birthdayDlq())
.to(dlxExchange())
.with(DLX_ROUTING_KEY_BIRTHDAY);
}
}Queue에서 메시지를 가져오는 리스너 (@RabbitListener는 Queue에 메시지가 있으면 상시적으로 가져옵니다):
@Service
@RequiredArgsConstructor
@Slf4j
public class WelcomeCouponMessageListener {
private final CouponService couponService;
@RabbitListener(queues = RabbitConfig.QUEUE_WELCOME)
public void receive(Long userId) {
log.info("RabbitMQ -> 회원가입 이벤트 수신. userId={}", userId);
try {
couponService.issueWelcomeCoupon(userId);
} catch (Exception e) {
log.error("웰컴 쿠폰 발급 실패 userId={}", userId, e);
throw e; // 예외 발생 시 DLX로 넘어감
}
}
}생일 쿠폰 스케줄러 + DLQ 재처리 스케줄러 예시(요지):
- 매월 1일 0시, 이번 달 생일자 목록을 1000건씩 페이징 조회해
coupon.birthday라우팅 키로 메시지 발행 - 매일 0시, DLQ에 쌓인 메시지를
receiveAndConvert로 하나씩 꺼내 원본 큐로 재발행(복구) — 큐가 비면 종료
자세한 전체 코드는 Notion 원본 참고.