apache/seatunnel · warning

MQTT source connection lost for client

Error message

MQTT source connection lost for client [{}], auto-reconnect will attempt recovery

What it means

The MQTT source reader's connectionLost callback fired after the broker dropped the connection. The reader records disconnectedSinceMs and disconnectCause, then logs this warning; Paho auto-reconnect is expected to restore the connection and resubscribe. If reconnect does not succeed within the configured timeout, the reader fails the task (see tests referencing reconnect timeout).

Solutions

  1. Confirm auto-reconnect restores the session; if the task eventually fails, check the reconnect timeout setting and increase it.
  2. Use a unique clientId per source reader instance to avoid broker kick loops.
  3. Inspect the logged disconnectCause (often reason code 32109 / socket loss) to fix the network or broker issue.
  4. Tune keepAlive interval to tolerate network latency; use ssl:// for unstable links.
  5. Verify broker-side connection limits aren't evicting clients.

Example fix

// before
source {
  Mqtt {
    client_id = "seatunnel-reader"
  }
}
// after
source {
  Mqtt {
    client_id = "seatunnel-reader-${subtask_index}"
  }
}
Defensive patterns

Strategy: try-catch

Try / catch

@Override
public void connectionLost(Throwable cause) {
    disconnectedSinceMs = System.currentTimeMillis();
    disconnectCause = cause;
    LOG.warn("MQTT connection lost, auto-reconnect in progress", cause);
}

Prevention

When it happens

Trigger: MqttSourceReader.connectionLost(Throwable) is invoked by Paho when the source client's connection drops; it logs with the configured clientId.

Common situations: Broker restart or crash; keep-alive timeout on idle networks; duplicate clientId causing broker kick; TLS/network issues; NAT idle timeouts.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/4c10b74ece3a0754. 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:151

                if (sourceConfig.isCleanSession()) {
                    mqttClient.unsubscribe(sourceConfig.getTopic());
                }
                mqttClient.disconnect();
            } else {
                mqttClient.disconnectForcibly();
            }
            mqttClient.close();
            LOG.info("MQTT source reader [{}] closed", sourceConfig.getClientId());
        } catch (MqttException e) {
            throw new IOException("Error closing MQTT source client", e);
        }
    }

    @Override
    public void connectionLost(Throwable cause) {
        disconnectedSinceMs = currentTimeMillis.getAsLong();
        disconnectCause = cause;
        LOG.warn(
                "MQTT source connection lost for client [{}], auto-reconnect will attempt recovery",
                sourceConfig.getClientId(),
                cause);
    }

    @Override
    public void connectComplete(boolean reconnect, String serverURI) {
        if (!reconnect) {
            return;
        }
        try {
            subscribeTopic();
            disconnectedSinceMs = -1L;
            disconnectCause = null;
            LOG.info(
                    "MQTT source reader [{}] resubscribed to topic [{}] after reconnect to [{}]",
                    sourceConfig.getClientId(),
                    sourceConfig.getTopic(),

View on GitHub (pinned to cf67b549a7)