Giải pháp xử lý tin nhắn trong RabbitMQ: Chết, trễ, đảm bảo, trùng lặp và quá tải

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ợpMô tả
Đơn hàng không thanh toánTin 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ầnSau nhiều lần thử → đưa vào hàng đợi chết để xử lý thủ công
Giảm tải đột ngộtHà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:

  1. TTL + Hàng đợi chết
  2. 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ợpThời gian trì hoãn
Đơn hàng tự động hủy15–30 phút
Xác thực SMS60 giây
Thông báo giao hàng24 giờ
Xử lý hàng loạt10 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ấtXác nhận + Lưu trữ
BrokerExchange, Queue, Message đều bền vững
ConsumptionACK 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ểmNhược điểmỨng dụng
Memory SetTốc độ caoMất khi khởi động, giới hạn bộ nhớĐơn máy, tải thấp
RedisTốc độ cao, bền vữngRedis hỏng thì mấtKhởi chạy cao
Cơ sở dữ liệuĐộ tin cậy caoTốc độ chậm hơnYê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ápMô tả
LimitingNgăn dòng dữ liệu bùng nổ
MonitoringGiám sát kích thước hàng đợi, cảnh báo khi vượt ngưỡng
IsolationTá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ăngGiải pháp chính
Hàng đợi chếtReject/timeout/đầy → Dead Letter Exchange → Queue chết
Hàng đợi trì hoãnTTL + chết hoặc plugin trì hoãn
Bảo đảm tin nhắnXác nhận sản xuất + 3 lớp lưu trữ + ACK thủ công
Đồng nhấtID 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

Thẻ: rabbitmq Dead Letter Queue Delayed Message Message Acknowledgment Idempotency

Đăng vào ngày 1 tháng 10 lúc 23:13