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

When WRITE_DISTRIBUTION_MODE=range is set but the sink has equality fields (primary keys), FlinkSink ignores the range request and instead hashes rows by the equality fields, logging this warning. Range distribution with primary keys is not always safe for the Flink streaming writer, so hash-by-key is used for backward compatibility instead of throwing.

Solutions

  1. Accept the fallback (hash by equality fields) — this is the safe behavior.
  2. Set write.distribution-mode=hash explicitly to make intent clear and silence the warning.
  3. Remove equalityFieldColumns if range distribution without keys is truly intended.

Example fix

// before
props.put("write.distribution-mode", "range"); // table has primary keys
// after
props.put("write.distribution-mode", "hash");
Defensive patterns

Strategy: validation

Validate before calling

if (mode == DistributionMode.RANGE && !eqFieldIds.isEmpty()) {
  mode = DistributionMode.HASH; // range is unsafe with primary keys in streaming writer
}

Prevention

When it happens

Trigger: distributeDataStream with DistributionMode.RANGE and non-empty equalityFieldIds — e.g. setting write.distribution-mode=range on a table with identifier fields / equalityFieldColumns configured.

Common situations: Tuning jobs for large backfill writes with range mode while forgetting the table has primary keys; upgrading a batch-style job config to a streaming sink.

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

Appendix: source

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