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
- 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)
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
- 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
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
- ILLEGAL_ARGUMENT
- OPERATION_NOT_SUPPORTED
- tables_configs[ ]: 'start_mode.timestamp' and…
- The read columns configuration will be filtered by the…
- All candidate sink tables were skipped during job parsing.
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)