NiFi ConsumeJMS与ActiveMQ Artemis持久化连接消息丢失问题排查
核心排查与解决方向
保证会话与订阅的状态一致性
尽管clientID和订阅名称未变,但每次重启处理器生成新会话名称,可能让Artemis判定为新订阅会话。需确保NiFi处理器重启时复用同一JMS会话实例,或创建会话时显式关联原有持久化订阅元数据。Artemis的持久化订阅绑定clientID+订阅名+会话组合,若会话重建未正确恢复订阅状态,服务器不会推送未确认消息。检查消息确认机制
确认处理器使用CLIENT_ACKNOWLEDGE或DUPS_OK_ACKNOWLEDGE模式,而非AUTO_ACKNOWLEDGE。自动确认模式下,处理器停机前未处理的消息可能被标记为已确认,重启后不会重新推送。需在代码中显式设置会话确认模式,且仅在FlowFile处理完成后调用message.acknowledge()。验证Artemis服务器订阅状态
用artemis queue stat命令或控制台查看目标队列的持久化订阅状态,确认对应clientID+订阅名的订阅存在且有未确认消息堆积。若订阅不存在,说明处理器重启未正确恢复订阅;若消息已被删除,可能是服务器的消息过期或订阅清理规则导致。规范处理器连接生命周期
确保处理器停止时正确关闭会话和连接,避免强制中断。异常关闭会导致Artemis未及时更新订阅的未确认消息状态,重启后服务器可能认为客户端仍在线,不会推送积压消息。在onStopped方法中添加session.close()和connection.close()逻辑,并捕获异常确保资源释放。确认订阅创建方式
必须使用createDurableSubscriber()(Topic模式)或createDurableConsumer()(Artemis扩展的Queue模式)创建持久化订阅,而非普通的createConsumer()。误用非持久化创建方法,即使设置clientID也无法实现持久化订阅。
关键代码调整示例
在NiFi处理器的会话创建逻辑中,固定订阅标识:
// 设置全局唯一且固定的clientID connection.setClientID("fixed-nifi-client-id"); // Topic场景创建持久化订阅 Topic topic = session.createTopic("target-topic"); MessageConsumer consumer = session.createDurableSubscriber(topic, "fixed-subscription-name"); // Queue场景使用Artemis专属持久化消费者 Queue queue = session.createQueue("target-queue"); MessageConsumer consumer = session.createDurableConsumer(queue, "fixed-subscription-name", null, false);
内容的提问来源于stack exchange,提问作者Akash Rai

