apache/seatunnel · error · SeaTunnelJsonFormatException
COMMON_UNSUPPORTED_DATA_TYPE (CommonErrorCodeDeprecated.UNSUPPORTED_DATA_TYPE)
COMMON_UNSUPPORTED_DATA_TYPE (CommonErrorCodeDeprecated.UNSUPPORTED_DATA_TYPE)
Error message
Unsupported format: %s
What it means
DefaultSeaTunnelRowSerializer.createSerializationSchema switches on the configured 'format' option and throws SeaTunnelJsonFormatException with code UNSUPPORTED_DATA_TYPE for any format outside the supported set (JSON, Text, CANAL_JSON, DEBEZIUM_JSON, AVRO, PROTOBUF, etc.). Despite the exception class name, it is a format-dispatch failure, not a JSON parsing failure.
Solutions
- Change the format option to a supported sink format: json, text, canal_json, debezium_json, avro, or protobuf.
- Check the connector documentation for the exact list of supported serialization formats and their exact spelling.
- If you need an unsupported format, add/enable a format plugin and extend the switch in createSerializationSchema.
Example fix
# before format = "csv" # not a supported Kafka sink format # after format = "json"
Defensive patterns
Strategy: validation
Validate before calling
java.util.Set<String> supported = java.util.Set.of(
"json", "text", "canal_json", "debezium_json", "avro", "protobuf");
String format = pluginConfig.get("format");
if (format == null || !supported.contains(format.toLowerCase())) {
throw new IllegalArgumentException("unsupported kafka sink format: " + format);
} Try / catch
try {
schema = DefaultSeaTunnelRowSerializer.createSerializationSchema(rowType, pluginConfig);
} catch (SeaTunnelJsonFormatException e) {
// log valid format values and rethrow with guidance
throw e;
} Prevention
- Copy format values from the official connector docs verbatim.
- Validate the HOCON config (format is an enum-like option) before job submission.
- Keep source/sink format lists separate — a source-supported format may not be sink-supported.
When it happens
Trigger: Setting the Kafka sink 'format' option to an unsupported value (e.g. 'csv', 'parquet', a misspelled name like 'jsn', or a format the connector has no case for), so the default branch executes.
Common situations: Typos in the format config value; using a format supported only by source (not sink) serialization; copy-pasting configs between connectors with different supported format lists.
Understand the failure class
Background: Invalid enum value errors: "Unknown type", "Invalid scope", "must be one of" — when a string is not on the library's allowed list — this error's family across 23 libraries.
Related errors
- COMMON_ILLEGAL_ARGUMENT (CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT)
- COMMON_ILLEGAL_ARGUMENT (CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT)
- COMMON_UNSUPPORTED_DATA_TYPE (CommonErrorCodeDeprecated.UNSUPPORTED_DATA_TYPE)
- start_mode.end_timestamp must not be greater than the…
- start_mode.timestamp must not be greater than the current…
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/63e3952790bd729b.
Report an issue: GitHub.
Appendix: source
Thrown at seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/serialize/DefaultSeaTunnelRowSerializer.java:438
case CANAL_JSON:
return new CanalJsonSerializationSchema(rowType);
case OGG_JSON:
return new OggJsonSerializationSchema(rowType);
case DEBEZIUM_JSON:
return new DebeziumJsonSerializationSchema(rowType);
case MAXWELL_JSON:
return new MaxWellJsonSerializationSchema(rowType);
case COMPATIBLE_DEBEZIUM_JSON:
return new CompatibleDebeziumJsonSerializationSchema(rowType, isKey);
case AVRO:
return new AvroSerializationSchema(rowType);
case PROTOBUF:
String protobufMessageName = pluginConfig.get(PROTOBUF_MESSAGE_NAME);
String protobufSchema = pluginConfig.get(PROTOBUF_SCHEMA);
return new ProtobufSerializationSchema(
rowType, protobufMessageName, protobufSchema);
default:
throw new SeaTunnelJsonFormatException(
CommonErrorCodeDeprecated.UNSUPPORTED_DATA_TYPE,
"Unsupported format: " + format);
}
}
private static Iterable<Header> convertToKafkaHeaders(Map<String, String> headersMap) {
if (MapUtils.isEmpty(headersMap)) {
return null;
}
RecordHeaders kafkaHeaders = new RecordHeaders();
for (Map.Entry<String, String> entry : headersMap.entrySet()) {
kafkaHeaders.add(
new RecordHeader(
entry.getKey(), entry.getValue().getBytes(StandardCharsets.UTF_8)));
}
return kafkaHeaders;
}
}View on GitHub (pinned to cf67b549a7)