apache/seatunnel · critical · MqttConnectorException

CONNECTION_FAILED

CONNECTION_FAILED

Error message

Failed to connect MQTT source client [${clientId}]

What it means

MqttSourceReader.open throws MqttConnectorException with code CONNECTION_FAILED when the underlying Paho client's connect or subscribe call throws MqttException. Before throwing, it quietly closes the client so no half-connected client leaks. Callers get the clientId in the message and the original MqttException as cause.

Solutions

  1. Read the chained MqttException reason code for the exact broker rejection cause
  2. Verify broker URL, port, credentials and ACL permissions in the source config
  3. Ensure client_id is unique across all running clients to avoid server kick-outs
  4. If using TLS, validate CA certificates and hostname settings; check network reachability (firewall, DNS)
  5. Let the engine retry the task if the outage is transient, or increase reconnect_timeout for later recovery

Example fix

// before
source {
  Mqtt {
    broker = "tcp://broker.local:1883"
    client_id = "seatunnel-source-1"
  }
}
// after
// verify broker reachable and credentials valid, e.g.:
// mosquitto_sub -h broker.local -p 1883 -t 'topic' -u user -P pass -i seatunnel-source-1
source {
  Mqtt {
    broker = "tcp://broker.local:1883"
    client_id = "seatunnel-source-1-unique"
    username = "user"
    password = "pass"
  }
}
Defensive patterns

Strategy: try-catch

Validate before calling

// preflight: verify broker connectivity before job submission
MqttClient probe = new MqttClient(brokerUrl, probeId, new MemoryPersistence());
probe.connect(new MqttConnectOptions());
probe.disconnect();

Type guard

null

Try / catch

try {
  reader.open(...);
} catch (MqttConnectorException e) {
  if (e.getCode() == MqttConnectorErrorCode.CONNECTION_FAILED
      && e.getCause() instanceof MqttException mEx) {
    LOG.error("MQTT connect failed, reason code {}", mEx.getReasonCode(), mEx);
  }
  throw e;
}

Prevention

When it happens

Trigger: MqttClient.connect() or subscribe() fails during reader open: broker unreachable, authentication rejected, TLS handshake failure, client_id already connected (server kick-out), or resubscribe after reconnect failing as in the referenced tests.

Common situations: Broker down or wrong broker URL/port; wrong username/password or ACLs denying subscription; duplicate client_id from another process or another SeaTunnel instance; TLS certificate mismatch.

Understand the failure class

Background: ECONNREFUSED and "connection refused" / "could not connect to server" errors: what they mean and how to fix them — this error's family across 44 libraries.

Related errors


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

    @Override
    public void open() {
        try {
            this.mqttClient =
                    new MqttClient(
                            sourceConfig.getUrl(),
                            sourceConfig.getClientId(),
                            new MemoryPersistence());
            this.mqttClient.setCallback(this);
            this.mqttClient.connect(buildConnectOptions(sourceConfig));
            subscribeTopic();
            LOG.info(
                    "MQTT source reader [{}] subscribed to topic [{}]",
                    sourceConfig.getClientId(),
                    sourceConfig.getTopic());
        } catch (MqttException e) {
            closeClientQuietly();
            throw new MqttConnectorException(
                    MqttConnectorErrorCode.CONNECTION_FAILED,
                    "Failed to connect MQTT source client [" + sourceConfig.getClientId() + "]",
                    e);
        }
    }

    @Override
    public void pollNext(Collector<SeaTunnelRow> output) throws Exception {
        checkReceiveException();
        checkReconnectTimeout();

        byte[] payload = messageQueue.poll(POLL_TIMEOUT_MS, TimeUnit.MILLISECONDS);
        if (payload == null) {
            return;
        }

        SeaTunnelRow row = deserializationSchema.deserialize(payload);
        if (row == null) {

View on GitHub (pinned to cf67b549a7)