{"record":{"id":"8a702c10465b4c57","repo":"apache/seatunnel","slug":"unsupported-mqtt-source-format-format-8a702c","errorCode":null,"errorMessage":"Unsupported MQTT source format: ${format}","messagePattern":"Unsupported MQTT source format: (.+?)","errorType":"exception","errorClass":"IllegalArgumentException","httpStatus":null,"severity":"error","filePath":"seatunnel-connectors-v2/connector-mqtt/src/main/java/org/apache/seatunnel/connectors/seatunnel/mqtt/source/MqttSourceReader.java","lineNumber":289,"sourceCode":"            options.setPassword(password.toCharArray());\n        }\n        return options;\n    }\n\n    private static DeserializationSchema<SeaTunnelRow> createDeserializationSchema(\n            MqttSourceConfig sourceConfig, CatalogTable catalogTable) {\n        SeaTunnelRowType rowType = catalogTable.getSeaTunnelRowType();\n        switch (sourceConfig.getFormat().toLowerCase()) {\n            case \"json\":\n                return new JsonDeserializationSchema(catalogTable, false, false);\n            case \"text\":\n                return TextDeserializationSchema.builder()\n                        .seaTunnelRowType(rowType)\n                        .delimiter(sourceConfig.getFieldDelimiter())\n                        .setCatalogTable(catalogTable)\n                        .build();\n            default:\n                throw new IllegalArgumentException(\n                        \"Unsupported MQTT source format: \" + sourceConfig.getFormat());\n        }\n    }\n\n    private void closeClientQuietly() {\n        if (mqttClient == null) {\n            return;\n        }\n        try {\n            mqttClient.close();\n        } catch (MqttException ignored) {\n            // Best-effort cleanup; the original connection exception is more important.\n        }\n    }\n}\n","sourceCodeStart":271,"sourceCodeEnd":305,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-mqtt/src/main/java/org/apache/seatunnel/connectors/seatunnel/mqtt/source/MqttSourceReader.java#L271-L305","documentation":"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.","triggerScenarios":"Configuring format = <typo or unsupported value> in the MQTT source (e.g. 'JSON' with wrong casing, 'csv', 'avro') and instantiating MqttSourceReader.","commonSituations":"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.","solutions":["Set format to a supported value: json, canonical_json, or text (check connector docs for exact list)","Fix casing — format matching is typically case-sensitive ('json', not 'JSON')","Remove the format option if you want the default","Consult the MQTT source documentation for the supported format enum"],"exampleFix":"// before\nformat = CSV\n// after\nformat = json","handlingStrategy":"validation","validationCode":"List<String> supported = List.of(\"json\", \"canonical_json\", \"text\");\nif (!supported.contains(config.getString(\"format\"))) {\n    throw new IllegalArgumentException(\n        \"format must be one of \" + supported + \", got: \" + config.getString(\"format\"));\n}","typeGuard":null,"tryCatchPattern":"try {\n    new MqttSourceReader(...);\n} catch (IllegalArgumentException e) {\n    LOG.error(\"Fix the 'format' option: {}\", e.getMessage());\n}","preventionTips":["Copy format values exactly from the MQTT source docs","Validate the job config with seatunnel.sh --check before running","Watch case sensitivity: 'json' not 'JSON'","Don't assume formats supported by other connectors work here"],"tags":["mqtt","config","unsupported-format","deserialization"],"backgroundTag":"unsupported-enum-value","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}