apache/seatunnel · error · KafkaConnectorException
COMMON_UNSUPPORTED_DATA_TYPE
COMMON_UNSUPPORTED_DATA_TYPE
Error message
Field name { %s } is not found! What it means
checkNativeSeaTunnelType() compares the input rowType against the schema required by SeaTunnel NATIVE format (nativeTableSchema().toPhysicalRowDataType()). If a required field cannot be found in the input rowType it throws UNSUPPORTED_DATA_TYPE. NATIVE format demands specific physical fields, not arbitrary user schemas.
Solutions
- Use a non-NATIVE format (JSON, TEXT, CANAL_JSON, etc.) for ordinary schemas
- Ensure the upstream data actually comes from a native-format producer matching nativeTableSchema()'s physical row type
- Rename/add the missing field so the input rowType contains the expected physical fields
Example fix
// before
Kafka { format = NATIVE } // input: arbitrary user table
// after
Kafka { format = JSON } Defensive patterns
Strategy: validation
Validate before calling
// choose NATIVE only for native-produced payloads
if (!isNativeProducedSchema(rowType)) {
config.putString("format", "JSON");
} Try / catch
try {
sink.open();
} catch (KafkaConnectorException e) {
if (e.getMessage().contains("is not found")) {
// fall back to JSON format and resubmit
}
} Prevention
- Use NATIVE format only when the upstream produces SeaTunnel native-format data
- Default to JSON/TEXT for arbitrary schemas
When it happens
Trigger: format = NATIVE with an input SeaTunnelRowType that lacks a field required by nativeTableSchema(); thrown from getSerializer() during writer init.
Common situations: Using NATIVE format without the corresponding upstream native source (e.g. SeaTunnelFile) that produces the expected schema; sending ordinary table data in NATIVE format.
Understand the failure class
Background: Schema validation failed / invalid input schema: payload rejected because its shape doesn't match the expected schema — this error's family across 28 libraries.
Related errors
- OPERATION_NOT_SUPPORTED
- COMMON-02
- COMMON-02
- COMMON_ILLEGAL_ARGUMENT (CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT)
- COMMON_ILLEGAL_ARGUMENT (CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT)
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/2d9eaebaa2ddc2b1.
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:318
CommonErrorCode.ILLEGAL_ARGUMENT,
String.format(
"Header field not found: %s, rowType: %s",
headerField, rowTypeFieldNames));
}
}
return headerFields;
}
return Collections.emptyList();
}
private void checkNativeSeaTunnelType(SeaTunnelRowType seaTunnelRowType) {
SeaTunnelRowType exceptRowType = nativeTableSchema().toPhysicalRowDataType();
for (int i = 0; i < exceptRowType.getFieldTypes().length; i++) {
String exceptField = exceptRowType.getFieldNames()[i];
SeaTunnelDataType<?> exceptFieldType = exceptRowType.getFieldTypes()[i];
int fieldIndex = seaTunnelRowType.indexOf(exceptField, false);
if (fieldIndex < 0) {
throw new KafkaConnectorException(
CommonErrorCode.UNSUPPORTED_DATA_TYPE,
String.format("Field name { %s } is not found!", exceptField));
}
SeaTunnelDataType<?> fieldType = seaTunnelRowType.getFieldType(fieldIndex);
if (exceptFieldType.getSqlType() != fieldType.getSqlType()) {
throw new KafkaConnectorException(
CommonErrorCode.UNSUPPORTED_DATA_TYPE,
String.format(
"Field name { %s } unsupported sql type { %s } !",
exceptField, fieldType.getSqlType()));
}
}
}
private TableSchema nativeTableSchema() {
return TableSchema.builder()
.column(
PhysicalColumn.of(View on GitHub (pinned to cf67b549a7)