apache/seatunnel · critical · IOException
MqttConnectorErrorCode.PUBLISH_FAILED
MqttConnectorErrorCode.PUBLISH_FAILED
Error message
Failed to publish MQTT message after
What it means
Thrown by publishWithRetry in MqttSinkWriter when an MQTT message still cannot be published after exhausting the retry window (retryTimeoutMs). It wraps MqttConnectorException with code MqttConnectorErrorCode.PUBLISH_FAILED and preserves the last broker-side exception as cause. This is the sink's terminal failure for an unflushable message.
Solutions
- Check the lastException cause chain for the root MqttException (reason code, socket error)
- Verify the broker is reachable and the connection options (broker URL, keep-alive) are correct in sink config
- Increase retry-timeout to tolerate longer broker outages, or add broker-side redundancy
- Check for client_id collisions if multiple writers share the same client_id
Example fix
// before
throw new IOException(
new MqttConnectorException(MqttConnectorErrorCode.PUBLISH_FAILED,
"Failed to publish MQTT message after " + retryTimeoutMs + "ms").getMessage(),
lastException);
// after
// keep the error, but prevent it: ensure broker availability and a unique client_id,
// or raise retryTimeoutMs so transient broker restarts are absorbed by the retry loop Defensive patterns
Strategy: retry
Validate before calling
// preflight broker reachability MqttClient c = new MqttClient(brokerUrl, clientId, new MemoryPersistence()); c.connect(opts); // fail fast before job submission
Type guard
null
Try / catch
try {
writer.flush();
} catch (IOException e) {
Throwable cause = e.getCause();
if (cause instanceof MqttConnectorException
&& ((MqttConnectorException) cause).getCode() == MqttConnectorErrorCode.PUBLISH_FAILED) {
// inspect cause.getCause() for root MqttException reason code
LOG.error("MQTT publish exhausted retries", cause);
}
throw e;
} Prevention
- Configure a retry-timeout that covers expected broker restart windows
- Ensure unique client_id per writer to avoid disconnect loops
- Add broker monitoring/alerting for availability before running long jobs
- Keep messages small to reduce broker timeout risk
When it happens
Trigger: publish() keeps throwing (connection loss, broker unavailable, client disconnected) for longer than retryTimeoutMs, so the retry loop exits and the final IOException is thrown from flushBuffer.
Common situations: MQTT broker down or restarting; network partition between SeaTunnel worker and broker; client_id conflict causing repeated disconnects; QoS 1/2 broker overload.
Understand the failure class
Background: "API request failed": what wrapped HTTP errors from external APIs mean and how to find the real cause — this error's family across 29 libraries.
Related errors
- Interrupted during MQTT publish retry
- batch_size must be >= 1, got:
- Can not sync pipeline owned slot profiles with IMap
- clean_session=false may cause broker-side state…
- client_id is required when clean_session=false for MQTT…
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/cc090d97b3600156.
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:211
MqttException lastException = null;
while (System.currentTimeMillis() < deadline) {
try {
if (mqttClient.isConnected()) {
mqttClient.publish(topic, message);
return;
}
} catch (MqttException e) {
lastException = e;
log.warn("Transient MQTT publish failure, retrying...", e);
}
try {
Thread.sleep(RETRY_BACKOFF_MS);
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
throw new IOException("Interrupted during MQTT publish retry", ie);
}
}
throw new IOException(
new MqttConnectorException(
MqttConnectorErrorCode.PUBLISH_FAILED,
"Failed to publish MQTT message after " + retryTimeoutMs + "ms")
.getMessage(),
lastException);
}
private static MqttConnectOptions buildConnectOptions(ReadonlyConfig config) {
MqttConnectOptions options = new MqttConnectOptions();
options.setAutomaticReconnect(true);
boolean cleanSession = config.get(MqttSinkOptions.CLEAN_SESSION);
options.setCleanSession(cleanSession);
if (!cleanSession) {
log.warn(
"clean_session=false may cause broker-side state accumulation. Ensure proper clientId management.");
}
options.setConnectionTimeout(config.get(MqttSinkOptions.CONNECTION_TIMEOUT));
View on GitHub (pinned to cf67b549a7)