apache/seatunnel · error · MqttConnectorException
RECEIVE_FAILED
RECEIVE_FAILED
Error message
Interrupted while buffering MQTT source message
What it means
In messageArrived (the Paho callback), the reader puts incoming messages into a bounded queue. If the thread waiting to enqueue is interrupted, it restores the interrupt flag, records the exception, and throws MqttConnectorException with RECEIVE_FAILED. Interruption here means the task is being cancelled or shutting down.
Solutions
- This usually signals intentional shutdown — verify whether a cancel/failover was in progress
- Increase the message queue capacity (source reader options) so messageArrived does not block on a full buffer
- Slow down upstream publish rate or increase downstream throughput to avoid the full-queue block
- Check the recorded receiveException surfaced later by pollNext (error 2122) for root cause
Example fix
// before
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new MqttConnectorException(...RECEIVE_FAILED, "Interrupted while buffering...", e);
}
// after (no code change needed; ensure task shutdown is expected and queue sizing is adequate)
Defensive patterns
Strategy: try-catch
Try / catch
try {
reader.pollNext(...);
} catch (MqttConnectorException e) {
if (e.getCause() instanceof InterruptedException) {
Thread.currentThread().interrupt(); // propagate shutdown
return;
}
throw e;
} Prevention
- Size the message queue for peak publish rate
- Ensure downstream pollNext keeps up (adequate parallelism)
- Avoid interrupting reader threads outside of intentional shutdown
- Always restore the interrupt flag after catching InterruptedException
When it happens
Trigger: The thread executing messageArrived is interrupted while blocked on queue.put (queue full) or queue.offer — i.e. task cancellation/stop arrives while the buffer is full.
Common situations: Zeta task failover or job cancel while the sink side is slow; queue capacity too small so the reader blocks on put during interruption; backpressure from downstream pollNext stalls.
Related errors
- Interrupted during MQTT publish retry
- batch_size must be >= 1, got:
- clean_session=false may cause broker-side state…
- client_id is required when clean_session=false for MQTT…
- CONNECTION_FAILED
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/cae6b9ac71d7c75f.
Report an issue: GitHub.
Appendix: source
Thrown at seatunnel-connectors-v2/connector-mqtt/src/main/java/org/apache/seatunnel/connectors/seatunnel/mqtt/source/MqttSourceReader.java:202
public void messageArrived(String topic, MqttMessage message) throws Exception {
if (message == null || message.getPayload() == null) {
return;
}
byte[] payload = Arrays.copyOf(message.getPayload(), message.getPayload().length);
try {
if (!messageQueue.offer(payload, QUEUE_OFFER_TIMEOUT_MS, TimeUnit.MILLISECONDS)) {
MqttConnectorException exception =
new MqttConnectorException(
MqttConnectorErrorCode.RECEIVE_FAILED,
"MQTT source message queue is full. Increase max_queue_size "
+ "or reduce MQTT message throughput.");
receiveException = exception;
throw exception;
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
receiveException = e;
throw new MqttConnectorException(
MqttConnectorErrorCode.RECEIVE_FAILED,
"Interrupted while buffering MQTT source message",
e);
}
}
@Override
public void deliveryComplete(IMqttDeliveryToken token) {
// Source-only client — outbound delivery acknowledgements are not expected.
}
void subscribeTopic() throws MqttException {
mqttClient.subscribe(sourceConfig.getTopic(), sourceConfig.getQos());
}
private void checkReceiveException() {
if (receiveException == null) {
return;View on GitHub (pinned to cf67b549a7)