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
- Read the chained MqttException reason code for the exact broker rejection cause
- Verify broker URL, port, credentials and ACL permissions in the source config
- Ensure client_id is unique across all running clients to avoid server kick-outs
- If using TLS, validate CA certificates and hostname settings; check network reachability (firewall, DNS)
- 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
- Verify broker URL, port, credentials and ACLs before deployment
- Guarantee client_id uniqueness across processes and clusters
- Validate TLS/CA configuration when using ssl:// brokers
- Monitor broker availability and set reconnect_timeout for transient outage recovery
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
- MQTT source connection lost for client
- Airtable API rate limit reached, retry
- At least one source plugin must be configured.
- batch_size must be >= 1, got:
- BigtableSourceSplitEnumerator already closed; cannot create…
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)