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
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.
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
Example fix
// before
.equalityFields("user_id", "ts")
// after (match table identifier fields)
.equalityFields(table.schema().identifierFieldNames().toArray(new String[0])) Defensive patterns
Strategy: validation
Validate before calling
Set<Integer> configured = SinkUtil.checkAndGetEqualityFieldIds(table, configuredColumns);
if (!configured.equals(table.schema().identifierFieldIds())) {
throw new IllegalArgumentException("Equality fields must match table identifier fields: " + table.schema().identifierFieldIds());
} Prevention
- 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
When it happens
Trigger: 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.
Common situations: 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.
Understand the failure class
Background: Schema validation failed / invalid input schema: payload rejected because its shape doesn't match the expected schema — this error's family across 28 libraries.
Related errors
- Hash distribute rows by equality fields, even though
- Hash distribute rows by equality fields, even though
- The configured equality field column IDs
- Can not alter the default database when the iceberg catalog…
- Cannot apply unknown modify-column change:
AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12).
Data as JSON: /api/errors/baffa0e0f76b02d1.
Report an issue: GitHub.
Appendix: source
Thrown at flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/sink/SinkUtil.java:74
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)