[ 
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)

Reply via email to