{"record":{"id":"e4f9eadf7e6cbd77","repo":"apache/iceberg","slug":"the-configured-equality-field-column-ids-are-no-e4f9ea","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/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/SinkUtil.java","lineNumber":72,"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":54,"sourceCodeEnd":90,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/SinkUtil.java#L54-L90","documentation":"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.","triggerScenarios":"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).","commonSituations":"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).","solutions":["Align the job's equality-field columns with the table schema's identifier fields (primary key).","Update the table schema identifier fields via schema evolution if the new equality columns are intended.","If the job-specified columns are intentional, acknowledge the warning — behavior already defaults to the job-specified columns."],"exampleFix":"// before\nnew FlinkSink.Builder().equalityFieldNames(Arrays.asList(\"order_id\", \"ts\"))\n// after (matching table PK)\nnew FlinkSink.Builder().equalityFieldNames(Collections.singletonList(\"order_id\"))","handlingStrategy":"validation","validationCode":"Set<String> configured = resolveEqualityFieldIds(table, configuredColumns);\nif (!configured.equals(table.schema().identifierFieldIds())) {\n  throw new IllegalArgumentException(\"Equality fields must match schema identifier fields: \"\n      + configured + \" vs \" + table.schema().identifierFieldIds());\n}","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Derive equality fields from table.schema().identifierFieldIds() instead of hardcoding.","Re-check config after primary-key schema evolution.","Keep equality column names in one shared config source."],"tags":["flink","iceberg-sink","equality-fields","primary-key"],"backgroundTag":"conflicting-config-options","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"}