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

  1. This usually signals intentional shutdown — verify whether a cancel/failover was in progress
  2. Increase the message queue capacity (source reader options) so messageArrived does not block on a full buffer
  3. Slow down upstream publish rate or increase downstream throughput to avoid the full-queue block
  4. 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

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


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)