apache/iceberg · warning

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, requesting range distribution while equality fields (primary keys) are configured triggers this warning: the sink ignores the range request and hashes rows by the equality fields instead, because range distribution with primary keys is not always safe in the Flink streaming writer. Kept as a fallback rather than an error for backward compatibility.

Solutions

  1. Set distributionMode to HASH explicitly so intent matches actual behavior.
  2. Remove equalityFieldColumns if range distribution without keys is intended.
  3. Leave as-is if the hash-by-key fallback is acceptable; it is safe and only logs.

Example fix

// before
IcebergSink.builder().distributionMode(DistributionMode.RANGE)... // equality fields set
// after
IcebergSink.builder().distributionMode(DistributionMode.HASH)...
Defensive patterns

Strategy: validation

Validate before calling

if (mode == DistributionMode.RANGE && !equalityFieldIds.isEmpty()) {
  mode = DistributionMode.HASH; // match actual sink behavior
}

Prevention

When it happens

Trigger: IcebergSink.builder().distributionMode(DistributionMode.RANGE) (or write.distribution-mode=range) with non-empty equalityFieldIds, in the range-distribution path of the builder.

Common situations: Batch tuning configs (range) reused for streaming jobs on tables with primary keys; expectation that range mode balances write parallelism while the table's identifier fields force hash-by-key.

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/4c878c1e1d276eae. Report an issue: GitHub.

Appendix: source

Thrown at flink/v2.2/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)