{"record":{"id":"baffa0e0f76b02d1","repo":"apache/iceberg","slug":"the-configured-equality-field-column-ids-are-no-baffa0","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.3/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.3/flink/src/main/java/org/apache/iceberg/flink/sink/SinkUtil.java#L56-L92","documentation":"A warning in SinkUtil when the user-specified equality field column names resolve to field IDs that differ from the table schema's identifierFieldIds (the declared primary key). The job's configured columns win and are used as equality fields; this alerts that the write key differs from the table's declared identifier.","triggerScenarios":"Calling IcebergSink/FlinkSink .equalityFields(\"a\",\"b\") (or setting the job-level equality-field config) with column names that do not exactly match the table's PRIMARY KEY / identifierFieldIds set.","commonSituations":"Typos or renamed columns in the equality field config; intentionally upserting on a different key than the table's identifier; schema evolution changed identifier fields after the job config was written.","solutions":["Align the configured equality fields with the table schema identifierFieldIds, or update the table's identifier fields via ALTER TABLE SET IDENTIFIER FIELDS","If the override is intentional, no action needed — the job columns are used by default","Validate the column names against table.schema().columns() before launching the job"],"exampleFix":"// before\n.equalityFields(\"user_id\", \"ts\")\n// after (match table identifier fields)\n.equalityFields(table.schema().identifierFieldNames().toArray(new String[0]))","handlingStrategy":"validation","validationCode":"Set<Integer> configured = SinkUtil.checkAndGetEqualityFieldIds(table, configuredColumns);\nif (!configured.equals(table.schema().identifierFieldIds())) {\n  throw new IllegalArgumentException(\"Equality fields must match table identifier fields: \" + table.schema().identifierFieldIds());\n}","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Derive equality fields from table.schema().identifierFieldNames() instead of hardcoding","Re-check configs after schema evolution","Validate column names against the schema before job launch"],"tags":["flink","iceberg","equality-fields","primary-key"],"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-23T08:17:48.524Z"}