[
https://issues.apache.org/jira/browse/NIFI-16421?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18123892#comment-18123892
]
Mateusz Chrzanowski commented on NIFI-16421:
--------------------------------------------
Hi [~turcsanyip], thanks for picking this up so quickly. I'm not planning to
open a PR - I'm not able to contribute code at the moment due to restrictions
on my side, so please keep the ticket assigned to yourself.
> ConsumeMQTT (MQTT v3) loses the in-flight message on reconnect with Resume
> Session: the broker's redelivery arrives before the consuming callback is
> installed
> --------------------------------------------------------------------------------------------------------------------------------------------------------------
>
> Key: NIFI-16421
> URL: https://issues.apache.org/jira/browse/NIFI-16421
> Project: Apache NiFi
> Issue Type: Bug
> Components: Extensions
> Affects Versions: 2.11.0, 2.12.0
> Environment: MiNiFi Java 2.11.0 in Podman on Ubuntu 24.04. Brokers on
> the same host: HiveMQ Edge 2025.13 and Eclipse Mosquitto 2. MQTT v3 AUTO
> (Paho), QoS 1, Session State = Resume Session.
> Reporter: Mateusz Chrzanowski
> Assignee: Peter Turcsanyi
> Priority: Major
> Labels: mqtt
>
> With {{Session State}} = _Resume Session_ and QoS 1, ConsumeMQTT is expected
> not to lose messages across a restart: the broker keeps the session and
> redelivers what was not acknowledged. In practice, *a graceful restart loses
> exactly one message* whenever a message is in flight at the moment the client
> disconnects. The message is logged at ERROR level and then acknowledged to
> the broker, so the broker never redelivers it again:
> {noformat}
> ERROR [MQTT Call: <client id>] o.a.nifi.processors.mqtt.ConsumeMQTT
> ConsumeMQTT[id=...]
> MQTT message arrived [topic:test/7/data; payload:[123, 34, 116, ...]]
> {noformat}
> h2. Cause
> In {{PahoMqttClientAdapter}} (nifi-mqtt-processors, {{rel/nifi-2.11.0}}):
> # The constructor calls {{client.setCallback(new DefaultMqttCallback())}}.
> Its {{messageArrived}} only logs the "MQTT message arrived" line shown above
> at ERROR and returns.
> # {{ConsumeMQTT.initializeClient()}} calls {{connect()}}, which connects with
> {{cleanSession = false}}.
> # It then calls {{subscribe()}}, which replaces the callback with {{new
> ConsumerMqttCallback(handler)}} and subscribes.
> MQTT 3.1.1 section 4.4 requires the server, when a client reconnects with
> CleanSession = 0, to re-send any unacknowledged PUBLISH packets. The broker
> does this immediately after CONNACK - between steps 2 and 3. The redelivered
> message reaches {{DefaultMqttCallback}}, which returns normally, and Paho
> then acknowledges it (no manual acks are used). The message is logged and
> discarded; the broker considers it delivered.
> h2. Steps to reproduce
> # An MQTT broker on the same host (reproduced with both HiveMQ Edge 2025.13
> and Mosquitto 2).
> # A flow: ConsumeMQTT (MQTT Specification Version {{v3 AUTO}}, QoS {{1}},
> Session State {{Resume Session}}, a fixed Client ID, topic filter
> {{test/+/data}}) -> PublishMQTT republishing every message unchanged to
> another topic.
> # A publisher sending a continuous stream (100 messages/s here), each message
> with a unique key, logging what it sent; a subscriber with a persistent
> session capturing the republished stream.
> # Restart MiNiFi gracefully several times while the stream runs ({{podman
> restart -t 60}}; the agent stops within about 3 s).
> # Compare what was published with what came out of the flow, key by key.
> h2. Expected
> No message lost across a graceful restart - QoS 1 plus a resumed session.
> h2. Actual - and the same runs with the fix below applied
> Five graceful restarts per run, 100 messages/s; every published message
> paired with its copy after the flow:
> ||Broker||NiFi MQTT bundle||Published||Lost||Lost per restart||ERROR "MQTT
> message arrived"||Duplicates||
> |HiveMQ Edge 2025.13|2.11.0 as released|37,254|*5*|1, 1, 1, 1, 1|5 - each one
> the lost message|0|
> |HiveMQ Edge 2025.13|2.11.0 + fix|37,254|*0*|0, 0, 0, 0, 0|0|0|
> |Mosquitto 2|2.11.0 as released|37,278|*5*|1, 1, 1, 1, 1|5 - each one the
> lost message|0|
> |Mosquitto 2|2.11.0 + fix|37,260|*0*|0, 0, 0, 0, 0|0|0|
> Every lost message was published within about 1 s after the restart was
> requested - the one in flight when the client disconnected - and matches the
> payload of its ERROR line. At a lower rate (48 messages/s, three restarts)
> two of three restarts lost one message and the third lost none: the loss
> needs a message in flight at the moment of disconnect, so its frequency
> depends on the rate.
> Expected to affect a processor stop/start as well (a new client is created
> and connected before {{subscribe()}}); only agent restarts were tested.
> The same code is in 2.12.0 and on main ({{PahoMqttClientAdapter}} and
> {{ConsumeMQTT.initializeClient()}} are unchanged as of 2026-10-01).
> h2. Suggested fix
> Install the consuming callback before {{connect()}}. The handler is already
> known to ConsumeMQTT at that point. A patch - prepared against
> {{rel/nifi-2.11.0}} and against main; three production files, eleven added
> lines - adds a default method to the {{MqttClient}} interface, implements it
> in the Paho adapter, and calls it in {{initializeClient()}} before
> {{connect()}}:
> {code:java}
> // MqttClient
> default void setReceivedMessageHandler(ReceivedMqttMessageHandler handler) {
> }
> // PahoMqttClientAdapter
> @Override
> public void setReceivedMessageHandler(ReceivedMqttMessageHandler handler) {
> client.setCallback(new ConsumerMqttCallback(handler));
> }
> // ConsumeMQTT.initializeClient()
> mqttClient = createMqttClient();
> mqttClient.setReceivedMessageHandler(this::handleReceivedMessage);
> mqttClient.connect();
> mqttClient.subscribe(topicPrefix + topicFilter, qos,
> this::handleReceivedMessage);
> {code}
> The default method is a no-op, so the v5 adapter and every other
> implementation are unchanged. The internal queue already exists before
> {{connect()}} (it is created in {{customValidate()}}), so a message delivered
> before {{subscribe()}} has somewhere to go.
> With it, a new unit test
> ({{TestConsumeMQTT.testQoS1NotCleanSessionMessageRedeliveredOnConnect}}; the
> test client redelivers a held message on {{connect()}}) fails without the
> {{ConsumeMQTT}} change and passes with it; all 41 tests of
> {{nifi-mqtt-processors}} pass and the {{contrib-check}} profile is clean. A
> pull request can follow this issue.
> A more complete alternative: {{setManualAcks(true)}} and acknowledge with
> {{messageArrivedComplete}} only after the message has been accepted. That
> would also close the loss on a hard kill, but it changes delivery semantics;
> the fix above does not.
> h2. Side observation
> {{ReceivedMqttMessage.isDuplicate()}} always returns {{false}} (also in
> 2.12.0: {{ConsumerMqttCallback}} does not pass the DUP flag on), so the
> {{mqtt.isDuplicate}} attribute and the {{_isDuplicate}} record field never
> reflect the MQTT DUP flag - a redelivered message cannot be told apart
> downstream. Separate from this issue; mentioned because it is what one would
> reach for to verify the fix.
> h2. Related
> * NIFI-10894 - _Resume Session in ConsumeMQTT does not work after NiFi
> restart_ (1.18/1.19, v5 adapter): the same class of defect - messages
> delivered before subscribe() are lost. Closed in January 2026 without a fix
> as part of the 1.x end-of-life cleanup, with a request to re-report for 2.x.
> * NIFI-10885 - ConsumeMQTT v5 client threads prevent a clean shutdown.
> ----
> _This report was drafted with the help of an AI assistant. The measurements,
> the code analysis and the patch come from my own lab runs and were checked by
> me._
--
This message was sent by Atlassian Jira
(v8.20.10#820010)