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

  1. Change the format option to a supported sink format: json, text, canal_json, debezium_json, avro, or protobuf.
  2. Check the connector documentation for the exact list of supported serialization formats and their exact spelling.
  3. 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

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


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)