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)