apache/seatunnel · error · KafkaConnectorException

OPERATION_NOT_SUPPORTED

OPERATION_NOT_SUPPORTED

Error message

kafka_headers_fields is not supported with NATIVE format. Please use JSON, TEXT, or other formats.

What it means

KafkaSinkWriter.getSerializer() rejects the kafka_headers_fields option when the message FORMAT is NATIVE, because native Kafka message-format serialization cannot attach custom headers derived from row fields. The library throws OPERATION_NOT_SUPPORTED to fail fast instead of silently dropping the headers.

Solutions

  1. Remove kafka_headers_fields from the sink config when using NATIVE format
  2. Switch format to JSON or TEXT if custom Kafka headers are required
  3. Put headers support behind a format that supports it and keep NATIVE only for pure payload serialization

Example fix

// before
Kafka {
  format = NATIVE
  kafka_headers_fields = ["hdr"]
}
// after
Kafka {
  format = JSON
  kafka_headers_fields = ["hdr"]
}
Defensive patterns

Strategy: validation

Validate before calling

// before submitting
boolean nativeFormat = "NATIVE".equalsIgnoreCase(config.getString("format"));
if (nativeFormat && config.getString("kafka_headers_fields") != null) {
  throw new IllegalArgumentException("kafka_headers_fields requires a non-NATIVE format");
}

Prevention

When it happens

Trigger: Configuring sink with format = NATIVE and a non-null kafka_headers_fields list; validation happens inside getSerializer() when the sink writer is created.

Common situations: User wants both SeaTunnel native serialization and custom Kafka headers; copy-pasting a config that worked with JSON/TEXT format and switching format to NATIVE.

Understand the failure class

Background: UnsupportedOperationException and "is not supported" errors: when a library deliberately refuses a call — this error's family across 30 libraries.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/fd74ab17fef19c7a. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaSinkWriter.java:188

        }

        kafkaProperties.put(
                ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, pluginConfig.get(BOOTSTRAP_SERVERS));
        kafkaProperties.put(
                ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class.getName());
        kafkaProperties.put(
                ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class.getName());
        return kafkaProperties;
    }

    private SeaTunnelRowSerializer<byte[], byte[]> getSerializer(
            ReadonlyConfig pluginConfig, SeaTunnelRowType seaTunnelRowType) {
        MessageFormat messageFormat = pluginConfig.get(FORMAT);
        String topic = pluginConfig.get(TOPIC);
        if (MessageFormat.NATIVE.equals(messageFormat)) {
            // Validate that kafka_headers_fields is not configured for NATIVE format
            if (pluginConfig.get(KAFKA_HEADERS_FIELDS) != null) {
                throw new KafkaConnectorException(
                        CommonErrorCode.OPERATION_NOT_SUPPORTED,
                        "kafka_headers_fields is not supported with NATIVE format. Please use JSON, TEXT, or other formats.");
            }
            checkNativeSeaTunnelType(seaTunnelRowType);
            return DefaultSeaTunnelRowSerializer.create(topic, messageFormat, seaTunnelRowType);
        }

        String delimiter = DEFAULT_FIELD_DELIMITER;

        if (pluginConfig.get(FIELD_DELIMITER) != null) {
            delimiter = pluginConfig.get(FIELD_DELIMITER);
        }
        if (pluginConfig.get(PARTITION_KEY_FIELDS) != null && pluginConfig.get(PARTITION) != null) {
            throw new KafkaConnectorException(
                    KafkaConnectorErrorCode.GET_TRANSACTIONMANAGER_FAILED,
                    "Cannot select both `partiton` and `partition_key_fields`. You can configure only one of them");
        }

View on GitHub (pinned to cf67b549a7)