{"record":{"id":"21f57d3cf9de93a1","repo":"apache/iceberg","slug":"the-configured-equality-field-column-ids-are-no-21f57d","errorCode":null,"errorMessage":"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.","messagePattern":"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\\.","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/sink/SinkUtil.java","lineNumber":74,"sourceCode":"\n  private static final Logger LOG = LoggerFactory.getLogger(SinkUtil.class);\n\n  static Set<Integer> checkAndGetEqualityFieldIds(Table table, List<String> equalityFieldColumns) {\n    Set<Integer> equalityFieldIds = Sets.newHashSet(table.schema().identifierFieldIds());\n    if (equalityFieldColumns != null && !equalityFieldColumns.isEmpty()) {\n      Set<Integer> equalityFieldSet = Sets.newHashSetWithExpectedSize(equalityFieldColumns.size());\n      for (String column : equalityFieldColumns) {\n        org.apache.iceberg.types.Types.NestedField field = table.schema().findField(column);\n        Preconditions.checkNotNull(\n            field,\n            \"Missing required equality field column '%s' in table schema %s\",\n            column,\n            table.schema());\n        equalityFieldSet.add(field.fieldId());\n      }\n\n      if (!equalityFieldSet.equals(table.schema().identifierFieldIds())) {\n        LOG.warn(\n            \"The configured equality field column IDs {} are not matched with the schema identifier field IDs\"\n                + \" {}, use job specified equality field columns as the equality fields by default.\",\n            equalityFieldSet,\n            table.schema().identifierFieldIds());\n      }\n      equalityFieldIds = Sets.newHashSet(equalityFieldSet);\n    }\n    return equalityFieldIds;\n  }\n\n  static long getMaxCommittedCheckpointId(\n      Table table, String flinkJobId, String operatorId, String branch) {\n    Snapshot snapshot = table.snapshot(branch);\n    long lastCommittedCheckpointId = INITIAL_CHECKPOINT_ID;\n\n    while (snapshot != null) {\n      Map<String, String> summary = snapshot.summary();\n      String snapshotFlinkJobId = summary.get(FLINK_JOB_ID);","sourceCodeStart":56,"sourceCodeEnd":92,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/sink/SinkUtil.java#L56-L92","documentation":"SinkUtil.checkAndGetEqualityFieldIds resolves the user-configured equality field columns to schema field IDs and compares them with the table schema's identifier field IDs. When the sets differ, it logs this warning and proceeds using the job-specified equality field columns. This matters for upsert mode: rows are matched on equality fields, which may silently diverge from the table's declared primary key (identifier fields).","triggerScenarios":"Calling SinkUtil.checkAndGetEqualityFieldIds (used when building Flink upsert sinks) with a table whose schema().identifierFieldIds() does not equal the set of IDs derived from the configured 'equality-field-columns' (e.g. via FlinkSink.Builder.setEqualityFieldColumns or table write.upsert.enabled with a partial column list).","commonSituations":"Table has an identifier (primary key) defined in Iceberg schema but the Flink job passes a different subset of columns as equality fields; identifiers were added to the table after the job was configured; typo in a configured column name so only some columns resolve.","solutions":["Make the configured equality field columns exactly match the table's identifier field columns.","Remove explicit equality-field-columns so the sink defaults to the schema identifier fields.","If the divergence is intentional, update the Iceberg schema identifiers (or accept the warning) and document that upsert matching uses the job-specified columns."],"exampleFix":"// before\nFlinkSink.forRowData(input)\n    .setEqualityFieldColumns(\"order_id\")\n    ...\n\n// after — match the table's identifier fields\nFlinkSink.forRowData(input)\n    .setEqualityFieldColumns(\"order_id\", \"line_number\")\n    ...","handlingStrategy":"validation","validationCode":"java.util.Set<Integer> configured = table.schema().columns().stream()\n    .filter(c -> equalityColumns.contains(c.name()))\n    .map(Types.NestedField::fieldId)\n    .collect(java.util.stream.Collectors.toSet());\nif (!configured.equals(table.schema().identifierFieldIds())) {\n  throw new IllegalArgumentException(\n      \"equality-field-columns must match schema identifier fields: \"\n          + table.schema().identifierFieldIds());\n}","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Derive equality columns from table.schema().identifierFieldIds() instead of hardcoding.","Re-check equality column config whenever the table's primary key changes.","Validate configured column names resolve to existing schema fields before submitting the job."],"tags":["flink","upsert","equality-fields","schema"],"backgroundTag":"schema-validation-failed","analyzedSha":"86d9c8fc543e7c56c9f624eb725f76c9baff9570","analyzedAt":"2026-09-12T00:46:39.097Z","contentChangedAt":"2026-09-12T00:46:39.097Z","schemaVersion":2},"datasetVersion":"2026-09-14T16:17:12.679Z"}