Khả năng lưu trữ bền vững trong ActiveMQ

Khái niệm về tin nhắn bền vững

Tin nhắn bền vững đảm bảo rằng mỗi tin nhắn chỉ được truyền đi một lần và được tiêu thụ thành công đúng một lần. Khi một tin nhắn được gửi đến đích với chế độ bền vững, dịch vụ nhắn tin sẽ lưu nó vào kho dữ liệu bền vững. Nếu hệ thống gặp sự cố (ví dụ: server bị sập), nó có thể khôi phục lại tin nhắn và tiếp tục chuyển đến người nhận. Mặc dù điều này làm tăng chi phí xử lý, nhưng đáng giá vì tính ổn định và độ tin cậy cao.

Có thể hiểu đơn giản: khi producer gửi tin nhắn thành công đến MQ, dù sau đó xảy ra lỗi như server chết, consumer mất kết nối, thì tin nhắn vẫn đảm bảo được tiêu thụ duy nhất một lần. Ngược lại, nếu tin nhắn chưa được gửi đến MQ, consumer sẽ không thể nhận được.

Chế độ bền vững và không bền vững cho hàng đợi (Queue)

1. Hàng đợi bền vững

Để thiết lập chế độ bền vững, sử dụng phương thức setDeliveryMode(DeliveryMode.PERSISTENT). Mặc định, các tin nhắn trong hàng đợi đều ở chế độ bền vững.

public class MessageProducer {
    public static final String BROKER_URL = "tcp://192.168.229.129:61616";
    public static final String USERNAME = "admin";
    public static final String PASSWORD = "admin";
    public static final String QUEUE_NAME = "orderQueue";
    public static final String MESSAGE_CONTENT = "Order-";

    public static void main(String[] args) throws JMSException {
        ActiveMQConnectionFactory factory = new ActiveMQConnectionFactory(USERNAME, PASSWORD, BROKER_URL);
        Connection connection = factory.createConnection();
        connection.start();

        Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        Queue queue = session.createQueue(QUEUE_NAME);
        MessageProducer sender = session.createProducer(queue);
        sender.setDeliveryMode(DeliveryMode.PERSISTENT);

        for (int i = 1; i <= 3; i++) {
            TextMessage msg = session.createTextMessage(MESSAGE_CONTENT + i);
            sender.send(queue, msg);
        }

        sender.close();
        session.close();
        connection.close();
        System.out.println("Đã gửi 3 tin nhắn đến hàng đợi.");
    }
}

2. Tiêu thụ tin nhắn

public class MessageConsumer {
    public static final String BROKER_URL = "tcp://192.168.229.129:61616";
    public static final String USERNAME = "admin";
    public static final String PASSWORD = "admin";
    public static final String QUEUE_NAME = "orderQueue";

    public static void main(String[] args) throws JMSException, IOException {
        ActiveMQConnectionFactory factory = new ActiveMQConnectionFactory(USERNAME, PASSWORD, BROKER_URL);
        Connection connection = factory.createConnection();
        connection.start();

        Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        Queue queue = session.createQueue(QUEUE_NAME);
        MessageConsumer receiver = session.createConsumer(queue);

        receiver.setMessageListener(message -> {
            if (message instanceof TextMessage textMsg) {
                try {
                    System.out.println("Tiêu thụ tin nhắn: " + textMsg.getText());
                } catch (JMSException e) {
                    e.printStackTrace();
                }
            }
        });

        System.in.read(); // Giữ chương trình chạy

        receiver.close();
        session.close();
        connection.close();
    }
}

Thực hiện các bước:

  1. Bắt đầu ActiveMQ.
  2. Chạy producer để gửi tin nhắn.
  3. Dừng ActiveMQ (giả lập lỗi máy chủ).
  4. Bắt lại ActiveMQ.
  5. Kiểm tra: tin nhắn vẫn tồn tại trong hàng đợi, consumer vẫn nhận được thông báo.

3. Hàng đợi không bền vững

Thay đổi dòng mã: producer.setDeliveryMode(DeliveryMode.NON_PERSISTENT);

Quá trình kiểm thử tương tự:

  1. Gửi tin nhắn.
  2. Dừng và khởi động lại ActiveMQ.
  3. Kết quả: tin nhắn bị mất hoàn toàn, consumer không nhận được gì.

Chế độ bền vững và không bền vững cho chủ đề (Topic)

1. Topic không bền vững (mặc định)

Topic mặc định hoạt động không bền vững. Nếu consumer không online khi producer gửi tin nhắn, thì tin nhắn sẽ bị bỏ qua — không thể thu hồi sau này.

2. Topic bền vững

Nếu consumer đã đăng ký trước (dùng setClientID() và createDurableSubscriber()), thì bất kể server có sập hay consumer offline, mọi tin nhắn gửi đến topic đều được lưu lại và giao cho consumer khi kết nối lại.

Producer cho topic bền vững

public class TopicPublisher {
    public static final String BROKER_URL = "tcp://192.168.229.129:61616";
    public static final String USERNAME = "admin";
    public static final String PASSWORD = "admin";
    public static final String TOPIC_NAME = "newsFeed";
    public static final String MSG_PREFIX = "Update-";

    public static void main(String[] args) throws JMSException {
        ActiveMQConnectionFactory factory = new ActiveMQConnectionFactory(USERNAME, PASSWORD, BROKER_URL);
        Connection connection = factory.createConnection();
        Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        Topic topic = session.createTopic(TOPIC_NAME);
        MessageProducer publisher = session.createProducer(topic);
        publisher.setDeliveryMode(DeliveryMode.PERSISTENT);

        connection.start();

        for (int i = 1; i <= 3; i++) {
            TextMessage msg = session.createTextMessage(MSG_PREFIX + i);
            publisher.send(msg);
        }

        publisher.close();
        session.close();
        connection.close();
        System.out.println("Gửi 3 tin nhắn tới topic.");
    }
}

Consumer cho topic bền vững

public class DurableSubscriber {
    public static final String BROKER_URL = "tcp://192.168.229.129:61616";
    public static final String USERNAME = "admin";
    public static final String PASSWORD = "admin";
    public static final String TOPIC_NAME = "newsFeed";

    public static void main(String[] args) throws JMSException, IOException {
        ActiveMQConnectionFactory factory = new ActiveMQConnectionFactory(USERNAME, PASSWORD, BROKER_URL);
        Connection connection = factory.createConnection();
        connection.setClientID("subscriber001"); // Đảm bảo client ID duy nhất

        Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
        Topic topic = session.createTopic(TOPIC_NAME);
        TopicSubscriber subscriber = session.createDurableSubscriber(topic, "newsGroup");

        connection.start();

        subscriber.setMessageListener(message -> {
            if (message instanceof TextMessage textMsg) {
                try {
                    System.out.println("Nhận được tin nhắn: " + textMsg.getText());
                } catch (JMSException e) {
                    e.printStackTrace();
                }
            }
        });

        System.in.read(); // Giữ vòng đời consumer

        subscriber.close();
        session.close();
        connection.close();
    }
}

Lưu ý quan trọng với Topic bền vững

  • Phải chạy consumer trước để đăng ký tên nhóm (subscription name).
  • Sau đó mới gửi tin nhắn từ producer.
  • Người tiêu dùng có thể offline mà không lo mất tin nhắn — khi kết nối lại, tất cả tin nhắn chưa đọc sẽ được gửi đến.

Thẻ: activemq JMS message persistence durable subscription topic

Đăng vào ngày 27 tháng 9 lúc 23:14