apache/seatunnel · error · KafkaConnectorException

COMMON_ILLEGAL_ARGUMENT (CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT)

COMMON_ILLEGAL_ARGUMENT (CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT)

Error message

Partition key field not found: %s, rowType: %s

What it means

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.

Solutions

  1. Fix the field name in partition_key_fields to exactly match a column of the input rowType (names in the error message)
  2. Check upstream Source/Transform output schema to confirm the column exists
  3. Correct the case of the field name — matching is case-sensitive (contains on exact strings)

Example fix

// before
partition_key_fields = ["UserID"]
// after (rowType has 'user_id')
partition_key_fields = ["user_id"]
Defensive patterns

Strategy: validation

Validate before calling

List<String> rowFields = Arrays.asList(seaTunnelRowType.getFieldNames());
partitionKeyFields.stream()
  .filter(f -> !rowFields.contains(f))
  .forEach(f -> { throw new IllegalArgumentException("Unknown partition key field: " + f); });

Prevention

When it happens

Trigger: partition_key_fields contains a column name that is not among seaTunnelRowType.getFieldNames(); thrown while resolving partition key fields before serialization.

Common situations: 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.

Understand the failure class

Background: 'Could not be found', 'does not exist', 'not found in database': the resource-not-found family when an ID, slug, key, or URI lookup comes back empty — this error's family across 20 libraries.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/3fdbe42736f24fb8. 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:279

        return transactionPrefix + "-" + checkpointId;
    }

    private void restoreState(List<KafkaSinkState> states) {
        if (!states.isEmpty()) {
            this.transactionPrefix = states.get(0).getTransactionIdPrefix();
            this.lastCheckpointId = states.get(0).getCheckpointId();
        }
    }

    private List<String> getPartitionKeyFields(
            ReadonlyConfig pluginConfig, SeaTunnelRowType seaTunnelRowType) {

        if (pluginConfig.get(PARTITION_KEY_FIELDS) != null) {
            List<String> partitionKeyFields = pluginConfig.get(PARTITION_KEY_FIELDS);
            List<String> rowTypeFieldNames = Arrays.asList(seaTunnelRowType.getFieldNames());
            for (String partitionKeyField : partitionKeyFields) {
                if (!rowTypeFieldNames.contains(partitionKeyField)) {
                    throw new KafkaConnectorException(
                            CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT,
                            String.format(
                                    "Partition key field not found: %s, rowType: %s",
                                    partitionKeyField, rowTypeFieldNames));
                }
            }
            return partitionKeyFields;
        }
        return Collections.emptyList();
    }

    private List<String> getHeaderFields(
            ReadonlyConfig pluginConfig, SeaTunnelRowType seaTunnelRowType) {

        if (pluginConfig.get(KAFKA_HEADERS_FIELDS) != null) {
            List<String> headerFields = pluginConfig.get(KAFKA_HEADERS_FIELDS);
            List<String> rowTypeFieldNames = Arrays.asList(seaTunnelRowType.getFieldNames());
            for (String headerField : headerFields) {

View on GitHub (pinned to cf67b549a7)