{"record":{"id":"2d9eaebaa2ddc2b1","repo":"apache/seatunnel","slug":"common-unsupported-data-type-2d9eae","errorCode":"COMMON_UNSUPPORTED_DATA_TYPE","errorMessage":"Field name { %s } is not found!","messagePattern":"Field name (.+?) is not found!","errorType":"error_code","errorClass":"KafkaConnectorException","httpStatus":null,"severity":"error","filePath":"seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaSinkWriter.java","lineNumber":318,"sourceCode":"                            CommonErrorCode.ILLEGAL_ARGUMENT,\n                            String.format(\n                                    \"Header field not found: %s, rowType: %s\",\n                                    headerField, rowTypeFieldNames));\n                }\n            }\n            return headerFields;\n        }\n        return Collections.emptyList();\n    }\n\n    private void checkNativeSeaTunnelType(SeaTunnelRowType seaTunnelRowType) {\n        SeaTunnelRowType exceptRowType = nativeTableSchema().toPhysicalRowDataType();\n        for (int i = 0; i < exceptRowType.getFieldTypes().length; i++) {\n            String exceptField = exceptRowType.getFieldNames()[i];\n            SeaTunnelDataType<?> exceptFieldType = exceptRowType.getFieldTypes()[i];\n            int fieldIndex = seaTunnelRowType.indexOf(exceptField, false);\n            if (fieldIndex < 0) {\n                throw new KafkaConnectorException(\n                        CommonErrorCode.UNSUPPORTED_DATA_TYPE,\n                        String.format(\"Field name { %s } is not found!\", exceptField));\n            }\n            SeaTunnelDataType<?> fieldType = seaTunnelRowType.getFieldType(fieldIndex);\n            if (exceptFieldType.getSqlType() != fieldType.getSqlType()) {\n                throw new KafkaConnectorException(\n                        CommonErrorCode.UNSUPPORTED_DATA_TYPE,\n                        String.format(\n                                \"Field name { %s } unsupported sql type { %s } !\",\n                                exceptField, fieldType.getSqlType()));\n            }\n        }\n    }\n\n    private TableSchema nativeTableSchema() {\n        return TableSchema.builder()\n                .column(\n                        PhysicalColumn.of(","sourceCodeStart":300,"sourceCodeEnd":336,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaSinkWriter.java#L300-L336","documentation":"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.","triggerScenarios":"format = NATIVE with an input SeaTunnelRowType that lacks a field required by nativeTableSchema(); thrown from getSerializer() during writer init.","commonSituations":"Using NATIVE format without the corresponding upstream native source (e.g. SeaTunnelFile) that produces the expected schema; sending ordinary table data in NATIVE format.","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"],"exampleFix":"// before\nKafka { format = NATIVE }  // input: arbitrary user table\n// after\nKafka { format = JSON }","handlingStrategy":"validation","validationCode":"// choose NATIVE only for native-produced payloads\nif (!isNativeProducedSchema(rowType)) {\n  config.putString(\"format\", \"JSON\");\n}","typeGuard":null,"tryCatchPattern":"try {\n  sink.open();\n} catch (KafkaConnectorException e) {\n  if (e.getMessage().contains(\"is not found\")) {\n    // fall back to JSON format and resubmit\n  }\n}","preventionTips":["Use NATIVE format only when the upstream produces SeaTunnel native-format data","Default to JSON/TEXT for arbitrary schemas"],"tags":["kafka","native-format","schema-validation","unsupported-data-type"],"backgroundTag":"schema-validation-failed","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}