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

  1. Set clean_session=true unless you specifically need broker-persisted state for a sink (rare — sinks publish, they don't consume).
  2. If clean_session=false is required, guarantee a unique clientId per writer instance (include subtask index).
  3. Clean up stale sessions on the broker (e.g. MQTT 5 session expiry interval, or broker admin cleanup).
  4. 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

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


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)