apache/seatunnel · warning

MQTT connection lost, auto-reconnect will attempt recovery

Error message

MQTT connection lost, auto-reconnect will attempt recovery

What it means

The MQTT sink's MqttCallback.connectionLost fired because the broker connection dropped. Automatic reconnect is enabled, so the writer deliberately logs a warning and does not throw; recovery is left to Paho's auto-reconnect. Buffered messages are flushed on the next successful connection.

Solutions

  1. No immediate action needed if automatic reconnect is enabled (default in this sink); watch for repeated occurrences.
  2. Ensure the clientId is unique per sink writer to avoid brokers kicking the connection.
  3. Tune keep-alive/connection timeout options to fit network conditions.
  4. Investigate the attached 'cause' for the root disconnect reason (e.g. 32109 socket loss).

Example fix

// before
url = "tcp://broker:1883"
// after (TLS for unstable networks)
url = "ssl://broker:8883"
Defensive patterns

Strategy: try-catch

Try / catch

// connectionLost is informational; detect stuck reconnection instead
@Override
public void connectionLost(Throwable cause) {
    log.warn("MQTT connection lost", cause);
    scheduleReconnectDeadlineCheck();
}

Prevention

When it happens

Trigger: The broker closes the socket or the network breaks while MqttSinkWriter is connected; Paho invokes connectionLost(Throwable).

Common situations: Broker restarts, keep-alive timeouts, LB idle connection eviction, network partitions, broker-side client kick due to duplicate clientId.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/04a195b478be083a. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-mqtt/src/main/java/org/apache/seatunnel/connectors/seatunnel/mqtt/sink/MqttSinkWriter.java:166

                try {
                    if (mqttClient.isConnected()) {
                        mqttClient.disconnect();
                    }
                    mqttClient.close();
                    log.info("MQTT sink writer closed");
                } catch (MqttException e) {
                    throw new IOException("Error closing MQTT client", e);
                }
            }
        }
    }

    // ---- MqttCallback implementation ----

    @Override
    public void connectionLost(Throwable cause) {
        // Auto-reconnect is enabled; log for observability but do not throw.
        log.warn("MQTT connection lost, auto-reconnect will attempt recovery", cause);
    }

    @Override
    public void messageArrived(String topic, MqttMessage message) {
        // Sink-only client — inbound messages are not expected.
    }

    @Override
    public void deliveryComplete(IMqttDeliveryToken token) {
        // QoS acknowledgement received from broker.
    }

    // ---- private helpers ----

    private void flushBuffer() throws IOException {
        if (messageBuffer.isEmpty()) {
            return;
        }

View on GitHub (pinned to cf67b549a7)