如何基于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
相关产品推荐
相关产品推荐

