{"record":{"id":"3fdbe42736f24fb8","repo":"apache/seatunnel","slug":"common-illegal-argument-commonerrorcodedeprecated-3fdbe4","errorCode":"COMMON_ILLEGAL_ARGUMENT (CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT)","errorMessage":"Partition key field not found: %s, rowType: %s","messagePattern":"Partition key field not found: (.+?), rowType: (.+?)","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":279,"sourceCode":"        return transactionPrefix + \"-\" + checkpointId;\n    }\n\n    private void restoreState(List<KafkaSinkState> states) {\n        if (!states.isEmpty()) {\n            this.transactionPrefix = states.get(0).getTransactionIdPrefix();\n            this.lastCheckpointId = states.get(0).getCheckpointId();\n        }\n    }\n\n    private List<String> getPartitionKeyFields(\n            ReadonlyConfig pluginConfig, SeaTunnelRowType seaTunnelRowType) {\n\n        if (pluginConfig.get(PARTITION_KEY_FIELDS) != null) {\n            List<String> partitionKeyFields = pluginConfig.get(PARTITION_KEY_FIELDS);\n            List<String> rowTypeFieldNames = Arrays.asList(seaTunnelRowType.getFieldNames());\n            for (String partitionKeyField : partitionKeyFields) {\n                if (!rowTypeFieldNames.contains(partitionKeyField)) {\n                    throw new KafkaConnectorException(\n                            CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT,\n                            String.format(\n                                    \"Partition key field not found: %s, rowType: %s\",\n                                    partitionKeyField, rowTypeFieldNames));\n                }\n            }\n            return partitionKeyFields;\n        }\n        return Collections.emptyList();\n    }\n\n    private List<String> getHeaderFields(\n            ReadonlyConfig pluginConfig, SeaTunnelRowType seaTunnelRowType) {\n\n        if (pluginConfig.get(KAFKA_HEADERS_FIELDS) != null) {\n            List<String> headerFields = pluginConfig.get(KAFKA_HEADERS_FIELDS);\n            List<String> rowTypeFieldNames = Arrays.asList(seaTunnelRowType.getFieldNames());\n            for (String headerField : headerFields) {","sourceCodeStart":261,"sourceCodeEnd":297,"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#L261-L297","documentation":"getPartitionKeyFields() verifies every name in partition_key_fields exists in the sink's SeaTunnelRowType; a missing name fails with ILLEGAL_ARGUMENT showing the field and the available rowType field names. This prevents runtime lookups against non-existent columns.","triggerScenarios":"partition_key_fields contains a column name that is not among seaTunnelRowType.getFieldNames(); thrown while resolving partition key fields before serialization.","commonSituations":"Typo in the field name; upstream transform renamed or dropped the column; case-sensitivity mismatch between config and schema; schema changed after the config was written.","solutions":["Fix the field name in partition_key_fields to exactly match a column of the input rowType (names in the error message)","Check upstream Source/Transform output schema to confirm the column exists","Correct the case of the field name — matching is case-sensitive (contains on exact strings)"],"exampleFix":"// before\npartition_key_fields = [\"UserID\"]\n// after (rowType has 'user_id')\npartition_key_fields = [\"user_id\"]","handlingStrategy":"validation","validationCode":"List<String> rowFields = Arrays.asList(seaTunnelRowType.getFieldNames());\npartitionKeyFields.stream()\n  .filter(f -> !rowFields.contains(f))\n  .forEach(f -> { throw new IllegalArgumentException(\"Unknown partition key field: \" + f); });","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Cross-check config field names against the upstream schema after any pipeline change","Match field name case exactly","Pin transform outputs with result_table_name and verify columns"],"tags":["kafka","config-validation","schema-mismatch","partition-key"],"backgroundTag":"resource-not-found","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"}