apache/seatunnel · warning
clean_session=false may cause broker-side state…
Error message
clean_session=false may cause broker-side state accumulation. Ensure proper clientId management.
What it means
When clean_session=false, the MQTT broker retains subscription and queued-message state per clientId across connections. The sink warns that without careful, unique clientId management, broker-side state (queued messages, session state) can accumulate indefinitely and consume broker resources.
Solutions
- Set clean_session=true unless you specifically need broker-persisted state for a sink (rare — sinks publish, they don't consume).
- If clean_session=false is required, guarantee a unique clientId per writer instance (include subtask index).
- Clean up stale sessions on the broker (e.g. MQTT 5 session expiry interval, or broker admin cleanup).
- Monitor broker session/storage metrics for growth.
Example fix
// before
sink {
Mqtt {
clean_session = false
}
}
// after
sink {
Mqtt {
clean_session = true
}
} Defensive patterns
Strategy: validation
Validate before calling
// validate sink config before submit
if (!config.getBoolean("clean_session")) {
log.warn("Persistent session requested: ensure unique clientId per writer");
} Prevention
- Default to clean_session=true for sinks
- If persistent sessions are needed, template clientIds with writer index
- Periodically purge stale broker sessions
- Watch broker retained-state/storage metrics
When it happens
Trigger: buildConnectOptions reads MqttSinkOptions.CLEAN_SESSION and it is false; the warning is logged during option construction before connecting.
Common situations: Users set clean_session=false expecting at-least-once semantics without realizing each writer instance needs a distinct persistent clientId; multiple parallel sink writers sharing a clientId cause session conflicts and state buildup.
Understand the failure class
Background: "Invalid value" and "allowed values are" config errors: what your library rejected and how to fix it — this error's family across 41 libraries.
Related errors
- client_id is required when clean_session=false for MQTT…
- MQTT connection lost, auto-reconnect will attempt recovery
- Airtable API rate limit reached, retry
- All candidate sink tables were skipped in Flink starter.
- All candidate sink tables were skipped in Flink starter.
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/9e354c037b0225da.
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:225
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));
String username = config.get(MqttSinkOptions.USERNAME);
if (username != null && !username.isEmpty()) {
options.setUserName(username);
}
String password = config.get(MqttSinkOptions.PASSWORD);
if (password != null && !password.isEmpty()) {
options.setPassword(password.toCharArray());
}
return options;
}
private static SerializationSchema createSerializationSchema(
SeaTunnelRowType rowType, ReadonlyConfig config) {
String format = config.get(MqttSinkOptions.FORMAT);View on GitHub (pinned to cf67b549a7)