Mateusz Chrzanowski created NIFI-16421:
------------------------------------------

             Summary: 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.12.0, 2.11.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


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