apache/iceberg · warning

The configured equality field column IDs

Error message

The configured equality field column IDs {} are not matched with the schema identifier field IDs {}, use job specified equality field columns as the equality fields by default.

What it means

SinkUtil validates that the equality field columns configured for the job resolve to exactly the schema's identifier field IDs. When they differ, it warns and uses the job-specified equality fields anyway. This flags a likely misconfiguration: the writer's dedup keys won't match the table's primary key definition.

Solutions

  1. Align the job's equality-field columns with the table schema's identifier fields (primary key).
  2. Update the table schema identifier fields via schema evolution if the new equality columns are intended.
  3. If the job-specified columns are intentional, acknowledge the warning — behavior already defaults to the job-specified columns.

Example fix

// before
new FlinkSink.Builder().equalityFieldNames(Arrays.asList("order_id", "ts"))
// after (matching table PK)
new FlinkSink.Builder().equalityFieldNames(Collections.singletonList("order_id"))
Defensive patterns

Strategy: validation

Validate before calling

Set<String> configured = resolveEqualityFieldIds(table, configuredColumns);
if (!configured.equals(table.schema().identifierFieldIds())) {
  throw new IllegalArgumentException("Equality fields must match schema identifier fields: "
      + configured + " vs " + table.schema().identifierFieldIds());
}

Prevention

When it happens

Trigger: Calling SinkUtil.checkAndGetEqualityFieldIds with a table whose schema identifierFieldIds differ from the columns given via the job's equality-field option (e.g. 'write.upsert.enabled' with configured equality columns that don't equal the table's primary key).

Common situations: Table primary key changed after the Flink job was written; typos or column renames in the equality-fields option; the table has no identifier fields but the job configures equality fields (or vice versa).

Understand the failure class

Background: Conflicting config options: "cannot be used together" — configuration validation errors across open-source libraries — this error's family across 162 libraries.

Related errors


AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12). Data as JSON: /api/errors/e4f9eadf7e6cbd77. Report an issue: GitHub.

Appendix: source

Thrown at flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/SinkUtil.java:72

  private static final Logger LOG = LoggerFactory.getLogger(SinkUtil.class);

  static Set<Integer> checkAndGetEqualityFieldIds(Table table, List<String> equalityFieldColumns) {
    Set<Integer> equalityFieldIds = Sets.newHashSet(table.schema().identifierFieldIds());
    if (equalityFieldColumns != null && !equalityFieldColumns.isEmpty()) {
      Set<Integer> equalityFieldSet = Sets.newHashSetWithExpectedSize(equalityFieldColumns.size());
      for (String column : equalityFieldColumns) {
        org.apache.iceberg.types.Types.NestedField field = table.schema().findField(column);
        Preconditions.checkNotNull(
            field,
            "Missing required equality field column '%s' in table schema %s",
            column,
            table.schema());
        equalityFieldSet.add(field.fieldId());
      }

      if (!equalityFieldSet.equals(table.schema().identifierFieldIds())) {
        LOG.warn(
            "The configured equality field column IDs {} are not matched with the schema identifier field IDs"
                + " {}, use job specified equality field columns as the equality fields by default.",
            equalityFieldSet,
            table.schema().identifierFieldIds());
      }
      equalityFieldIds = Sets.newHashSet(equalityFieldSet);
    }
    return equalityFieldIds;
  }

  static long getMaxCommittedCheckpointId(
      Table table, String flinkJobId, String operatorId, String branch) {
    Snapshot snapshot = table.snapshot(branch);
    long lastCommittedCheckpointId = INITIAL_CHECKPOINT_ID;

    while (snapshot != null) {
      Map<String, String> summary = snapshot.summary();
      String snapshotFlinkJobId = summary.get(FLINK_JOB_ID);

View on GitHub (pinned to 86d9c8fc54)