Khi hệ thống yêu cầu xử lý các giao dịch quan trọng như cập nhật số dư tài khoản và trạng thái thanh toán, việc mất tin nhắn gửi đến dịch vụ đơn hàng có thể dẫn đến trạng thái đơn hàng sai (ví dụ: chưa thanh toán nhưng thực tế đã thanh toán). Để giải quyết vấn đề này, RabbitMQ cung cấp cơ chế xác nhận từ phía Producer bao gồm:
- ID toàn cục duy nhất cho mỗi tin nhắn.
ConfirmCallback– xử lý kết quả khi Exchange nhận tin.ReturnCallback– xử lý khi Exchange không route được tin nhắn tới Queue.
Nguyên lý hoạt động
- Khi tin nhắn đến được Exchange nhưng không route tới Queue (do routingKey sai, queue chưa tồn tại hoặc chưa bind), RabbitMQ vẫn trả về ACK cho Producer vì Exchange đã nhận thành công. Tuy nhiên, đây là lỗi do cấu hình, không phải lỗi hệ thống, việc gửi lại có thể vô ích; thường chỉ cần ghi log.
- Khi Exchange không route được, RabbitMQ gọi
ReturnCallbackvới thông tin: tin nhắn, Exchange, routingKey, mã lỗi và mô tả lỗi. Điều này cho phép ghi log hoặc gửi lại tin nhắn. ConfirmCallbackđược gọi sau khi Exchange nhận tin nhắn. Nếu nhận thành công, trả vềack = true; nếu thất bại, trả vềnack = truekèm lý do. Producer dựa vào đó để quyết định gửi lại.
Cài đặt ReturnCallback
Một lớp cấu hình Spring Boot đăng ký callback này:
package com.example.config;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.ReturnedMessage;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.context.annotation.Configuration;
import javax.annotation.PostConstruct;
@Slf4j
@Configuration
public class RabbitConfig {
private final RabbitTemplate template;
public RabbitConfig(RabbitTemplate template) {
this.template = template;
}
@PostConstruct
public void init() {
template.setReturnsCallback(returned -> {
log.warn("Return callback triggered");
log.debug("Exchange: {}", returned.getExchange());
log.debug("Routing key: {}", returned.getRoutingKey());
log.debug("Message: {}", returned.getMessage());
log.debug("Reply code: {}", returned.getReplyCode());
log.debug("Reply text: {}", returned.getReplyText());
});
}
}
Cài đặt ConfirmCallback với CorrelationData
Mỗi tin nhắn được gán một CorrelationData chứa ID duy nhất, sau đó đăng ký callback bất đồng bộ:
@Test
void testConfirm() {
// 1. Tạo CorrelationData với UUID
CorrelationData corrData = new CorrelationData(java.util.UUID.randomUUID().toString());
// 2. Thêm callback cho Future của CorrelationData
corrData.getFuture().addCallback(new ListenableFutureCallback<CorrelationData.Confirm>() {
@Override
public void onFailure(Throwable ex) {
// Trường hợp lỗi Future (rất hiếm xảy ra)
log.error("Future failed", ex);
}
@Override
public void onSuccess(CorrelationData.Confirm result) {
if (result.isAck()) {
log.info("Message ACK – Exchange nhận thành công");
} else {
log.error("Message NACK – lý do: {}", result.getReason());
}
}
});
// 3. Gửi tin nhắn với CorrelationData
template.convertAndSend("exchange.demo", "invalid.key", "Hello", corrData);
}
Test các tình huống
- Gửi thành công: Tin nhắn đến Exchange và route đến Queue – nhận ACK, không có Return.
- RoutingKey sai: Tin nhắn đến Exchange, không có Queue phù hợp – vẫn nhận ACK (vì Exchange nhận OK), đồng thời
ReturnCallbackđược kích hoạt.
Lưu ý: Cơ chế xác nhận khiến Producer phải chờ phản hồi từ RabbitMQ, ảnh hưởng hiệu suất. Vì tỉ lệ lỗi rất thấp, không nên bật nếu không thực sự cần.
Kết nối lại (Reconnection) từ Producer
Trong trường hợp mất kết nối tới RabbitMQ, Producer có thể được cấu hình tự động thử lại. Mặc định tính năng này tắt.
spring:
rabbitmq:
template:
retry:
enabled: true
initial-interval: 1000ms # thời gian chờ lần đầu
multiplier: 2 # mỗi lần sau nhân lên 2 lần
max-attempts: 3
Ví dụ: lần đầu chờ 1 giây, lần thứ hai chờ 2 giây, lần thứ ba chờ 4 giây.
Hạn chế: trong khi chờ kết nối lại, thread hiện tại bị chặn, các tác vụ khác không thể thực thi. Vì vậy cần cấu hình thời gian chờ hợp lý (khoảng 200ms) để tránh ảnh hưởng luồng xử lý.