apache/iceberg · error · ValidationException

Invalid primary key ' '. A primary key must not contain…

Error message

Invalid primary key '%s'. A primary key must not contain duplicate columns. Found: %s

What it means

Validation copied from Flink's DefaultSchemaResolver: the primary key constraint lists one or more columns twice. The message names the key and the duplicate columns found by comparing the PK column list against the schema's column-name lookup — a malformed table definition, not a data issue.

Solutions

  1. Deduplicate the primary key column list before declaring the constraint
  2. Fix the DDL so each column appears once in PRIMARY KEY clause
  3. Ensure programmatic schema construction doesn't append the same key column twice

Example fix

// before
PRIMARY KEY (id, id) NOT ENFORCED
// after
PRIMARY KEY (id) NOT ENFORCED
Defensive patterns

Strategy: validation

Validate before calling

Set<String> pk = new LinkedHashSet<>(primaryKey.getColumns());
if (pk.size() != primaryKey.getColumns().size()) {
  throw new IllegalArgumentException("duplicate primary key columns");
}

Try / catch

try { FlinkSchemaUtil.toResolvedSchema(schema); } catch (ValidationException e) { /* fix PK definition */ }

Prevention

When it happens

Trigger: Defining PRIMARY KEY (a, a, b) or a duplicated column arising from programmatic ResolvedSchema construction, before converting to an Iceberg schema via toResolvedSchema.

Common situations: Generated SQL/DDL building primary keys from a collection with duplicate entries; schema evolving tooling concatenating key columns.

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/iceberg@86d9c8fc54 (2026-09-12). Data as JSON: /api/errors/f922c80da367e6fb. Report an issue: GitHub.

Appendix: source

Thrown at flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/FlinkSchemaUtil.java:349

    return new ResolvedSchema(columns, Collections.emptyList(), uniqueConstraint);
  }

  /**
   * Copied from
   * org.apache.flink.table.catalog.DefaultSchemaResolver#validatePrimaryKey(org.apache.flink.table.catalog.UniqueConstraint,
   * java.util.List)
   */
  private static void validatePrimaryKey(UniqueConstraint primaryKey, List<Column> columns) {
    final Map<String, Column> columnsByNameLookup =
        columns.stream().collect(Collectors.toMap(Column::getName, Function.identity()));

    final Set<String> duplicateColumns =
        primaryKey.getColumns().stream()
            .filter(name -> Collections.frequency(primaryKey.getColumns(), name) > 1)
            .collect(Collectors.toSet());

    if (!duplicateColumns.isEmpty()) {
      throw new ValidationException(
          String.format(
              "Invalid primary key '%s'. A primary key must not contain duplicate columns. Found: %s",
              primaryKey.getName(), duplicateColumns));
    }

    for (String columnName : primaryKey.getColumns()) {
      Column column = columnsByNameLookup.get(columnName);
      if (column == null) {
        throw new ValidationException(
            String.format(
                "Invalid primary key '%s'. Column '%s' does not exist.",
                primaryKey.getName(), columnName));
      }

      if (!column.isPhysical()) {
        throw new ValidationException(
            String.format(
                "Invalid primary key '%s'. Column '%s' is not a physical column.",

View on GitHub (pinned to 86d9c8fc54)