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

  1. Use a non-NATIVE format (JSON, TEXT, CANAL_JSON, etc.) for ordinary schemas
  2. Ensure the upstream data actually comes from a native-format producer matching nativeTableSchema()'s physical row type
  3. 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

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


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)