You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何基于Quarkus实现MQTT消息可靠处理,避免丢失?

问题描述

使用Quarkus 3.9.4、SmallRye连接器及ActiveMQ Artemis 2.33开发Java应用,核心功能是接收MQTT消息,处理后转发至AMQP队列。为保证消息全程不丢失(允许重复),已配置MQTT QoS 1及Clean Session=false。

当前遇到的问题:

  • 网络离线恢复后,代理能正常推送积压消息,符合预期;
  • 但应用在处理大量积压消息时异常停止,SmallRye本地队列中未处理的消息会丢失;
  • 尝试使用@Acknowledgment(Acknowledgment.Strategy.MANUAL)注解实现手动确认,但未确认的消息代理不会重发。
解决方案

1. 修正手动确认的配置与代码实现

配置调整

在application.properties中明确开启SmallRye MQTT连接器的手动确认模式:

mp.messaging.incoming.mqtt-connector.acknowledgment-mode=manual
mp.messaging.incoming.mqtt-connector.qos=1
mp.messaging.incoming.mqtt-connector.clean-session=false

代码实现

在消息消费方法中,正确获取Message对象,仅在处理成功后调用ack()确认;处理失败时调用nack()触发代理重发:

import io.smallrye.reactive.messaging.mqtt.MqttMessageMetadata;
import org.eclipse.microprofile.reactive.messaging.Incoming;
import org.eclipse.microprofile.reactive.messaging.Message;
import org.eclipse.microprofile.reactive.messaging.Acknowledgment;
import org.jboss.logging.Logger;

import javax.enterprise.context.ApplicationScoped;
import java.util.concurrent.CompletionStage;

@ApplicationScoped
public class MqttMessageHandler {
    private static final Logger LOG = Logger.getLogger(MqttMessageHandler.class);

    @Incoming("mqtt-connector")
    @Acknowledgment(Acknowledgment.Strategy.MANUAL)
    public CompletionStage<Void> processMessage(Message<String> message) {
        String messageId = message.getMetadata(MqttMessageMetadata.class)
                                  .map(MqttMessageMetadata::getId)
                                  .orElse("unknown");
        
        try {
            // 执行业务处理与转发AMQP队列的逻辑
            handleAndForward(message.getPayload());
            LOG.info("消息处理完成,确认消息ID: " + messageId);
            return message.ack();
        } catch (Exception e) {
            LOG.error("消息处理失败,触发重发,消息ID: " + messageId, e);
            return message.nack(e);
        }
    }

    private void handleAndForward(String payload) {
        // 替换为实际业务处理及AMQP发送逻辑
    }
}

2. 禁用SmallRye预取机制

默认情况下,SmallRye会预取一批消息到本地队列,应用崩溃时这些预取未处理的消息无法被代理收回。调整预取数量为0,让消息始终保留在代理端直到确认:

mp.messaging.incoming.mqtt-connector.prefetch=0

3. 配置ActiveMQ Artemis的持久化与重发策略

确保Artemis开启消息持久化,并配置合理的重发规则,在broker.xml中添加对应地址的设置:

<address-setting match="mqtt.#">
    <persist-messages>true</persist-messages>
    <redelivery-delay>3000</redelivery-delay> <!-- 重发间隔3秒 -->
    <max-redelivery-attempts>-1</max-redelivery-attempts> <!-- 无限重发,可根据需求调整次数 -->
    <dead-letter-address>DLQ</dead-letter-address> <!-- 最终失败消息进入死信队列 -->
</address-setting>

4. 添加消息轨迹日志

通过日志记录消息的接收、处理、确认状态,方便排查问题:

// 在processMessage方法中添加日志(已包含在步骤1的代码示例中)

内容的提问来源于stack exchange,提问作者Ricky Habegger

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 12:52:29