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

  1. Accept the hash-by-equality-fields fallback, or switch the table property to write.distribution-mode=hash explicitly
  2. If range distribution is required, write with a batch writer or without equality fields
  3. 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

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


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)