apache/seatunnel · error · IllegalArgumentException

Unsupported MQTT source format

Error message

Unsupported MQTT source format: ${format}

What it means

createDeserializationSchema maps the configured 'format' option to a SeaTunnel DeserializationSchema. An unknown format falls into the default branch and throws IllegalArgumentException. It is a configuration validation failure: the format value is not one of json, canonical_json, text, etc.

Solutions

  1. Set format to a supported value: json, canonical_json, or text (check connector docs for exact list)
  2. Fix casing — format matching is typically case-sensitive ('json', not 'JSON')
  3. Remove the format option if you want the default
  4. Consult the MQTT source documentation for the supported format enum

Example fix

// before
format = CSV
// after
format = json
Defensive patterns

Strategy: validation

Validate before calling

List<String> supported = List.of("json", "canonical_json", "text");
if (!supported.contains(config.getString("format"))) {
    throw new IllegalArgumentException(
        "format must be one of " + supported + ", got: " + config.getString("format"));
}

Try / catch

try {
    new MqttSourceReader(...);
} catch (IllegalArgumentException e) {
    LOG.error("Fix the 'format' option: {}", e.getMessage());
}

Prevention

When it happens

Trigger: Configuring format = <typo or unsupported value> in the MQTT source (e.g. 'JSON' with wrong casing, 'csv', 'avro') and instantiating MqttSourceReader.

Common situations: Typos or case mismatches in the HOCON config; copying a format name from another connector that supports more formats; upgrading SeaTunnel and using a format this connector never supported.

Related errors


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

            options.setPassword(password.toCharArray());
        }
        return options;
    }

    private static DeserializationSchema<SeaTunnelRow> createDeserializationSchema(
            MqttSourceConfig sourceConfig, CatalogTable catalogTable) {
        SeaTunnelRowType rowType = catalogTable.getSeaTunnelRowType();
        switch (sourceConfig.getFormat().toLowerCase()) {
            case "json":
                return new JsonDeserializationSchema(catalogTable, false, false);
            case "text":
                return TextDeserializationSchema.builder()
                        .seaTunnelRowType(rowType)
                        .delimiter(sourceConfig.getFieldDelimiter())
                        .setCatalogTable(catalogTable)
                        .build();
            default:
                throw new IllegalArgumentException(
                        "Unsupported MQTT source format: " + sourceConfig.getFormat());
        }
    }

    private void closeClientQuietly() {
        if (mqttClient == null) {
            return;
        }
        try {
            mqttClient.close();
        } catch (MqttException ignored) {
            // Best-effort cleanup; the original connection exception is more important.
        }
    }
}

View on GitHub (pinned to cf67b549a7)