apache/iceberg · info
Hash distribute rows by equality fields, even though
Error message
Hash distribute rows by equality fields, even though {}=range is set. Range distribution for primary keys are not always safe in Flink streaming writer. What it means
In IcebergSink's range-distribution path, if equality fields (primary key) are set, range distribution is unsafe in a Flink streaming writer. The sink logs this warning and falls back to hash distribution keyed by equality fields for backward compatibility instead of throwing.
Solutions
- Accept the hash-by-equality-fields fallback, or switch the table property to write.distribution-mode=hash explicitly
- If range distribution is required, write with a batch writer or without equality fields
- Document that range mode is ignored for primary-key streaming writes
Example fix
// before
UPDATE table SET TBLPROPERTIES ('write.distribution-mode'='range');
// after
UPDATE table SET TBLPROPERTIES ('write.distribution-mode'='hash'); Defensive patterns
Strategy: validation
Validate before calling
if ("range".equals(table.properties().get(TableProperties.WRITE_DISTRIBUTION_MODE))
&& !equalityFieldIds.isEmpty()) {
// streaming + PK: range is unsafe; switch to hash before building the sink
table.updateProperties().set(TableProperties.WRITE_DISTRIBUTION_MODE, "hash").commit();
} Prevention
- Audit table properties when reusing batch tables for streaming writes
- Use hash mode for all primary-key Flink streaming sinks
- Log expected vs actual distribution mode at job startup
When it happens
Trigger: IcebergSink constructor's distribution logic hits WRITE_DISTRIBUTION_MODE=range (or a sortOrder implying range) with non-empty equalityFieldIds in a streaming write.
Common situations: Table property write.distribution-mode=range set for batch workloads and reused by a Flink streaming upsert job; users expecting global ordering but silently receiving hash-by-PK shuffles.
Understand the failure class
Background: Conflicting config options: "cannot be used together" — configuration validation errors across open-source libraries — this error's family across 162 libraries.
Related errors
- Hash distribute rows by equality fields, even though
- Fallback to use 'none' distribution mode, because there are…
- Fallback to use 'none' distribution mode, because there are…
- The configured equality field column IDs
- Can not alter the default database when the iceberg catalog…
AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12).
Data as JSON: /api/errors/9e29c0e102c62e98.
Report an issue: GitHub.
Appendix: source
Thrown at flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergSink.java:1050
return Optional.ofNullable(flinkWriteConf.writeParallelism()).orElseGet(input::getParallelism);
}
private DataStream<RowData> distributeDataStreamByRangeDistributionMode(
DataStream<RowData> input,
Schema iSchema,
PartitionSpec partitionSpec,
SortOrder sortOrderParam) {
int writerParallelism = resolveWriterParallelism(input);
// needed because of checkStyle not allowing us to change the value of an argument
SortOrder sortOrder = sortOrderParam;
// Ideally, exception should be thrown in the combination of range distribution and
// equality fields. Primary key case should use hash distribution mode.
// Keep the current behavior of falling back to keyBy for backward compatibility.
if (!equalityFieldIds.isEmpty()) {
LOG.warn(
"Hash distribute rows by equality fields, even though {}=range is set. "
+ "Range distribution for primary keys are not always safe in "
+ "Flink streaming writer.",
WRITE_DISTRIBUTION_MODE);
return input.keyBy(new EqualityFieldKeySelector(iSchema, flinkRowType, equalityFieldIds));
}
// range distribute by partition key or sort key if table has an SortOrder
Preconditions.checkState(
sortOrder.isSorted() || partitionSpec.isPartitioned(),
"Invalid write distribution mode: range. Need to define sort order or partition spec.");
if (sortOrder.isUnsorted()) {
sortOrder = Partitioning.sortOrderFor(partitionSpec);
LOG.info("Construct sort order from partition spec");
}
LOG.info("Range distribute rows by sort order: {}", sortOrder);
StatisticsOrRecordTypeInformation statisticsOrRecordTypeInformation =View on GitHub (pinned to 86d9c8fc54)