1. Hàng đợi tin nhắn chết (Dead Letter Queue)
Tin nhắn chết là những tin nhắn không thể được xử lý bình thường.
- Chưa được xác nhận (reject/nack) và không quay lại hàng đợi.
- Hết thời gian sống (TTL).
- Hàng đợi đạt giới hạn độ dài tối đa.
Cấu hình hàng đợi chết
// Thiết lập tham số cho hàng đợi chính
Map<String, Object> args = new HashMap<>();
args.put("x-dead-letter-exchange", "dlx.exchange"); // Exchange xử lý tin nhắn chết
args.put("x-dead-letter-routing-key", "dlx.key"); // Routing key của tin nhắn chết
Queue normalQueue = new Queue("normal.queue", true, false, false, args);
// Khai báo exchange và queue chết
DirectExchange dlxExchange = new DirectExchange("dlx.exchange");
Queue dlq = new Queue("dlq", true);
// Gắn kết
BindingBuilder.bind(dlq).to(dlxExchange).with("dlx.key");
Thông tin bổ sung trong tin nhắn chết
MessageProperties props = message.getMessageProperties();
props.getDeadLetterExchange(); // Exchange xử lý chết
props.getDeadLetterRoutingKey(); // Routing key chết
props.getExpiration(); // TTL nếu do hết hạn
props.getRedelivered(); // Đã được gửi lại chưa
Ứng dụng thực tế
| Trường hợp | Mô tả |
| Đơn hàng không thanh toán | Tin nhắn hết hạn → vào hàng đợi chết → hủy đơn |
| Xử lý thất bại nhiều lần | Sau nhiều lần thử → đưa vào hàng đợi chết để xử lý thủ công |
| Giảm tải đột ngột | Hàng đợi đầy → tin nhắn vượt giới hạn → chuyển sang hàng đợi chết |
2. Hàng đợi trì hoãn (Delay Queue)
Nguyên lý
Tin nhắn được gửi nhưng không xử lý ngay, mà chờ một khoảng thời gian trước khi được tiêu thụ.
Phương pháp triển khai:
- TTL + Hàng đợi chết
- Sử dụng plugin:
rabbitmq_delayed_message_exchange
Phương án 1: TTL + Hàng đợi chết
// Cấu hình hàng đợi trì hoãn
Map<String, Object> args = new HashMap<>();
args.put("x-message-ttl", 5000); // Chờ 5 giây
args.put("x-dead-letter-exchange", "order.exchange");
args.put("x-dead-letter-routing-key", "order.created");
Queue delayQueue = new Queue("delay.queue", false, false, false, args);
// Gán exchange và binding
channel.exchangeDeclare("delay.exchange", "direct", true);
channel.queueDeclare("delay.queue", false, false, false, args);
channel.queueBind("delay.queue", "delay.exchange", "delay.key");
// Exchange và queue xử lý thật sự
channel.exchangeDeclare("order.exchange", "direct", true);
channel.queueDeclare("order.queue", true, false, false, null);
channel.queueBind("order.queue", "order.exchange", "order.created");
Phương án 2: Plugin trì hoãn (ưu tiên)
# Kích hoạt plugin
rabbitmq-plugins enable rabbitmq_delayed_message_exchange
// Tạo exchange trì hoãn
Map<String, Object> args = new HashMap<>();
args.put("x-delayed-type", "direct");
channel.exchangeDeclare(
"delay.exchange",
"x-delayed-message",
true,
false,
args
);
// Gửi tin nhắn với độ trễ
AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder()
.deliveryMode(2)
.headers(Map.of("x-delay", 5000)) // Trì hoãn 5 giây
.build();
channel.basicPublish("delay.exchange", "delay.key", properties, body);
Ứng dụng phổ biến
| Trường hợp | Thời gian trì hoãn |
| Đơn hàng tự động hủy | 15–30 phút |
| Xác thực SMS | 60 giây |
| Thông báo giao hàng | 24 giờ |
| Xử lý hàng loạt | 10 phút |
3. Bảo đảm tin nhắn đáng tin cậy
Nguyên nhân mất tin nhắn
- Sản xuất thất bại trước khi đến Broker.
- Broker lỗi hoặc mất dữ liệu.
- Consumer xử lý sai hoặc crash.
Phía sản xuất (Producer)
// Bật xác nhận
channel.confirmSelect();
channel.addConfirmListener(
(ack, deliveryTag) -> {
// Tin nhắn đã đến Broker thành công
},
(ack, deliveryTag) -> {
// Thất bại → retry
}
);
// Gán ID duy nhất
String correlationId = UUID.randomUUID().toString();
AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder()
.correlationId(correlationId)
.deliveryMode(2) // Duy trì trên disk
.build();
channel.basicPublish(exchange, routingKey, properties, message);
Phía Broker
// Xác nhận bền vững
channel.exchangeDeclare(exchange, "direct", true); // durable
channel.queueDeclare(queue, true, false, false, null); // durable
// Tin nhắn lưu trữ
AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder()
.deliveryMode(2)
.build();
Phía tiêu thụ (Consumer)
channel.basicConsume(queue, false, new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope,
AMQP.BasicProperties properties, byte[] body) {
try {
process(body);
channel.basicAck(envelope.getDeliveryTag(), false); // ACK thủ công
} catch (Exception e) {
channel.basicNack(envelope.getDeliveryTag(), false, true); // Từ chối, có thể quay lại
}
}
});
Tóm tắt giải pháp bảo đảm
| Vị trí | Giải pháp |
| Sản xuất | Xác nhận + Lưu trữ |
| Broker | Exchange, Queue, Message đều bền vững |
| Consumption | ACK thủ công + xử lý trùng lặp |
4. Đảm bảo tính đồng nhất (Idempotency)
Nguyên nhân tin nhắn bị lặp
- Producer gửi lại do timeout.
- Consumer không ACK kịp.
- Network retransmission.
Giải pháp
Phương án 1: Dùng ID tin nhắn (Set in-memory)
String messageId = UUID.randomUUID().toString();
AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder()
.messageId(messageId)
.build();
private final Set<String> processedIds = new HashSet<>();
public void handleDelivery(Message message) {
String id = message.getMessageId();
if (processedIds.contains(id)) return;
process(message);
processedIds.add(id);
}
Phương án 2: Redis kiểm tra duy nhất
String key = "mq:dedup:" + messageId;
Boolean result = redisTemplate.opsForValue().setIfAbsent(key, "1", 24, TimeUnit.HOURS);
if (!result) return; // Đã xử lý
process(message);
Phương án 3: Kiểm tra bằng khóa chính cơ sở dữ liệu
CREATE TABLE mq_message_log (
message_id VARCHAR(64) PRIMARY KEY,
status VARCHAR(20),
create_time DATETIME
);
-- Insert khi xử lý
INSERT INTO mq_message_log (message_id, status) VALUES ('msg_123', 'PROCESSED');
-- Nếu lỗi (duplicate key) → bỏ qua
So sánh giải pháp
| Phương án | Ưu điểm | Nhược điểm | Ứng dụng |
| Memory Set | Tốc độ cao | Mất khi khởi động, giới hạn bộ nhớ | Đơn máy, tải thấp |
| Redis | Tốc độ cao, bền vững | Redis hỏng thì mất | Khởi chạy cao |
| Cơ sở dữ liệu | Độ tin cậy cao | Tốc độ chậm hơn | Yêu cầu độ tin cậy cực cao |
5. Xử lý tin nhắn bị dồn ứ (Backlog)
Nguyên nhân dồn ứ
- Tốc độ sản xuất > tốc độ tiêu thụ.
- Consumer lỗi hoặc khởi động lại.
- Logic xử lý phức tạp, tốn thời gian.
- Doanh thu tăng đột biến.
Giải pháp
Phương án 1: Tăng số lượng consumer
for (int i = 0; i < 10; i++) {
new Thread(() -> {
channel.basicConsume(queue, false, consumer);
}).start();
}
Phương án 2: Tăng khả năng xử lý tạm thời
// Tăng số lượng tin nhắn lấy về mỗi lần
channel.basicQos(100); // Lấy tối đa 100 tin nhắn
Phương án 3: Xử lý tin nhắn hết hạn
// Đặt TTL cho hàng đợi
Map<String, Object> args = new HashMap<>();
args.put("x-message-ttl", 3600000); // 1 giờ
args.put("x-dead-letter-exchange", "dlx.exchange");
// Hoặc đặt riêng từng tin nhắn
AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder()
.expiration("3600000")
.build();
Phương án 4: Phân luồng tin nhắn
if (isCritical) {
channel.basicPublish("critical.exchange", "key", properties, body);
} else {
log.warn("Bỏ qua tin nhắn không quan trọng: {}", messageId);
// Có thể ném vào dead letter hoặc xóa luôn
}
Giám sát và cảnh báo
Map<String, Object> queueInfo = channel.queueDeclarePassive(queue).getQueueArgs();
Long msgCount = (Long) queueInfo.get("messageCount");
if (msgCount > 10000) {
alert("Cảnh báo dồn ứ: " + queue + " có " + msgCount + " tin nhắn");
}
Biện pháp phòng ngừa
| Biện pháp | Mô tả |
| Limiting | Ngăn dòng dữ liệu bùng nổ |
| Monitoring | Giám sát kích thước hàng đợi, cảnh báo khi vượt ngưỡng |
| Isolation | Tách biệt tin nhắn quan trọng và không quan trọng |
| Dead letter | Đặt hàng đợi chết để xử lý tin nhắn thất bại |
6. Ví dụ thực tế: Hủy đơn hàng sau thời gian chờ
@Configuration
public class OrderDelayConfig {
@Bean
public DirectExchange orderDelayExchange() {
return new DirectExchange("order.delay.exchange", true, false);
}
@Bean
public Queue orderDelayQueue() {
Map<String, Object> args = new HashMap<>();
args.put("x-message-ttl", 1800000); // 30 phút
args.put("x-dead-letter-exchange", "order.dlx.exchange");
args.put("x-dead-letter-routing-key", "order.cancel");
return new Queue("order.delay.queue", false, false, false, args);
}
@Bean
public DirectExchange orderDlxExchange() {
return new DirectExchange("order.dlx.exchange", true, false);
}
@Bean
public Queue orderCancelQueue() {
return new Queue("order.cancel.queue", true);
}
@Bean
public Binding orderDelayBinding() {
return BindingBuilder.bind(orderDelayQueue())
.to(orderDelayExchange())
.with("order.delay");
}
@Bean
public Binding orderCancelBinding() {
return BindingBuilder.bind(orderCancelQueue())
.to(orderDlxExchange())
.with("order.cancel");
}
}
@Service
public class OrderService {
@Autowired
private RabbitTemplate rabbitTemplate;
public void createOrder(Order order) {
// Tạo đơn hàng...
rabbitTemplate.convertAndSend(
"order.delay.exchange",
"order.delay",
order.getOrderId(),
message -> {
message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
return message;
}
);
}
}
@Service
public class OrderCancelConsumer {
@RabbitListener(queues = "order.cancel.queue")
public void handleCancelOrder(String orderId) {
Order order = orderMapper.selectById(orderId);
if (order != null && "Pending".equals(order.getStatus())) {
order.setStatus("Cancelled");
orderMapper.updateById(order);
stockService.releaseStock(order.getProductId(), order.getQuantity());
log.info("Hủy đơn hàng do hết hạn: {}", orderId);
}
}
}
7. Tổng kết
| Tính năng | Giải pháp chính |
| Hàng đợi chết | Reject/timeout/đầy → Dead Letter Exchange → Queue chết |
| Hàng đợi trì hoãn | TTL + chết hoặc plugin trì hoãn |
| Bảo đảm tin nhắn | Xác nhận sản xuất + 3 lớp lưu trữ + ACK thủ công |
| Đồng nhất | ID tin nhắn + Redis/CSDL để tránh lặp |
| Dồn ứ | Tăng consumer + mở rộng tạm thời + loại bỏ theo TTL |