{"record":{"id":"40141e20fd006d5b","repo":"apache/iceberg","slug":"the-configured-equality-field-column-ids-are-no-40141e","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/FlinkSink.java","lineNumber":512,"sourceCode":"\n    @VisibleForTesting\n    List<Integer> checkAndGetEqualityFieldIds() {\n      List<Integer> equalityFieldIds = Lists.newArrayList(table.schema().identifierFieldIds());\n      if (equalityFieldColumns != null && !equalityFieldColumns.isEmpty()) {\n        Set<Integer> equalityFieldSet =\n            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 = Lists.newArrayList(equalityFieldSet);\n      }\n      return equalityFieldIds;\n    }\n\n    private DataStreamSink<Void> appendDummySink(SingleOutputStreamOperator<Void> committerStream) {\n      DataStreamSink<Void> resultStream =\n          committerStream\n              .sinkTo(new DiscardingSink<>())\n              .name(operatorName(String.format(\"IcebergSink %s\", this.table.name())))\n              .setParallelism(1);\n      if (uidPrefix != null) {\n        resultStream = resultStream.uid(uidPrefix + \"-dummysink\");","sourceCodeStart":494,"sourceCodeEnd":530,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/sink/FlinkSink.java#L494-L530","documentation":"Warning from FlinkSink.checkAndGetEqualityFieldIds when the equality-field columns configured on the job do not exactly match the table schema's identifier field IDs. The sink proceeds using the job-specified columns as equality fields, which can differ from what upsert semantics the table's identifier fields imply.","triggerScenarios":"Writing with FlinkSink in upsert mode while setEqualityFieldColumns picks columns whose resolved field IDs differ from table.schema().identifierFieldIds() — e.g. different columns, a subset, or the table's identifier fields changed after the job was configured.","commonSituations":"Table schema identifier fields updated (ALTER TABLE ... SET IDENTIFIER FIELDS) while the Flink job still passes old equality columns; typo'd or reordered column names; job config and table evolution drift.","solutions":["Align setEqualityFieldColumns(...) with the table's identifier fields, or remove it so the schema identifier fields are used","Verify current identifier field IDs with DESCRIBE TABLE / table.schema().identifierFieldIds()","If the job-specified columns are intentionally different, silence the warning by acknowledging the divergence in job docs/config review","Update and redeploy the Flink job after table identifier-field changes"],"exampleFix":"// before\nFlinkSink.forRowData(input)\n    .setEqualityFieldColumns(\"order_id\")\n    ...\n// after (match schema identifier fields, e.g. [order_id, line_number])\nFlinkSink.forRowData(input)\n    .setEqualityFieldColumns(\"order_id\", \"line_number\")\n    ...","handlingStrategy":"validation","validationCode":"Set<Integer> configured = equalityFieldSet;\nSet<Integer> identifiers = table.schema().identifierFieldIds();\nif (!identifiers.isEmpty() && !configured.equals(identifiers)) {\n  throw new IllegalArgumentException(\"equality columns must match identifier fields: \" + identifiers);\n}","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Before upsert writes, compare job equality columns with schema.identifierFieldIds()","Re-validate config after any ALTER TABLE SET/DROP IDENTIFIER FIELDS","Prefer omitting setEqualityFieldColumns when the schema identifier fields are correct","Add a startup assertion in the Flink job for this invariant"],"tags":["flink","upsert","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"}