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

When WRITE_DISTRIBUTION_MODE=range is configured but equality fields (primary key) are set, FlinkSink cannot safely apply range distribution in a streaming writer. It logs this warning and falls back to hash-distributing rows by the equality fields, preserving old behavior for backward compatibility.

Solutions

  1. Accept the fallback (hash by equality fields) - this is the safe behavior for primary-key streams
  2. Change the table's write.distribution-mode to hash for streaming primary-key writes
  3. If range distribution is truly needed for an append stream, clear equality fields and use the batch writer path

Example fix

// before
table.updateProperties().set(TableProperties.WRITE_DISTRIBUTION_MODE, "range").commit();
// after
table.updateProperties().set(TableProperties.WRITE_DISTRIBUTION_MODE, "hash").commit();
Defensive patterns

Strategy: validation

Validate before calling

if ("range".equals(table.properties().get(TableProperties.WRITE_DISTRIBUTION_MODE))
    && !equalityFieldIds.isEmpty()) {
  throw new IllegalStateException("range distribution unsupported with primary keys in Flink streaming; use hash");
}

Prevention

When it happens

Trigger: distributeDataStream encounters case RANGE with a non-empty equalityFieldIds while writing to a table with a primary key in streaming mode.

Common situations: Table property write.distribution-mode=range (often set for batch/Spark writes) reused by a Flink streaming upsert job; users expecting range shuffle ordering but getting hash by PK.

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

Appendix: source

Thrown at flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/sink/FlinkSink.java:667

              for (PartitionField partitionField : partitionSpec.fields()) {
                Preconditions.checkState(
                    equalityFieldIds.contains(partitionField.sourceId()),
                    "In 'hash' distribution mode with equality fields set, source column '%s' of partition field '%s' "
                        + "should be included in equality fields: '%s'",
                    table.schema().findColumnName(partitionField.sourceId()),
                    partitionField,
                    equalityFieldColumns);
              }
              return input.keyBy(new PartitionKeySelector(partitionSpec, iSchema, flinkRowType));
            }
          }

        case RANGE:
          // 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);

View on GitHub (pinned to 86d9c8fc54)