apache/iceberg · warning

: Fallback to use 'none' distribution mode, because there…

Error message

{}: Fallback to use 'none' distribution mode, because there are no equality fields set and {}='range' is not supported yet in flink

What it means

The dynamic sink's HashKeyGenerator does not support RANGE distribution in Flink. When a table requests RANGE with no equality fields, it warns and falls back to a table-level key selector ('none'-like behavior per table).

Solutions

  1. Set write.distribution-mode to 'hash' or 'none' on tables written by the Flink dynamic sink.
  2. Provide equality fields if hash-keyed distribution is needed alongside keyed semantics.
  3. Ignore the warning if per-table serialization is acceptable for this workload.

Example fix

// before
table.updateProperties().set("write.distribution-mode", "range").commit();
// after
table.updateProperties().set("write.distribution-mode", "none").commit();
Defensive patterns

Strategy: fallback

Validate before calling

if ("range".equals(table.properties().get("write.distribution-mode"))) {
  LOG.warn("Range distribution unsupported in Flink dynamic sink; falling back to per-table keying");
}

Prevention

When it happens

Trigger: A DynamicRecord for a table whose write.distribution-mode is RANGE, with empty equality fields, reaching HashKeyGenerator.getKeySelector.

Common situations: Tables configured for batch/range distribution in Spark also being written by the Flink dynamic sink; default table properties set range mode.

Related errors


AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12). Data as JSON: /api/errors/7230a7fe66984c99. Report an issue: GitHub.

Appendix: source

Thrown at flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/sink/dynamic/HashKeyGenerator.java:177

              Types.NestedField sourceField = schema.findField(partitionField.sourceId());
              Preconditions.checkState(
                  sourceField != null && equalityFields.contains(sourceField.name()),
                  "%s: In 'hash' distribution mode with equality fields set, partition field '%s' "
                      + "should be included in equality fields: '%s'",
                  tableName,
                  partitionField,
                  schema.columns().stream()
                      .filter(c -> equalityFields.contains(c.name()))
                      .collect(Collectors.toList()));
            }
            return partitionKeySelector(
                tableName, schema, spec, writeParallelism, maxWriteParallelism);
          }
        }

      case RANGE:
        if (equalityFields.isEmpty()) {
          LOG.warn(
              "{}: Fallback to use 'none' distribution mode, because there are no equality fields set "
                  + "and {}='range' is not supported yet in flink",
              tableName,
              WRITE_DISTRIBUTION_MODE);
          return tableKeySelector(tableName, writeParallelism, maxWriteParallelism);
        } else {
          LOG.info(
              "{}: Distribute rows by equality fields, because there are equality fields set "
                  + "and {}='range' is not supported yet in flink",
              tableName,
              WRITE_DISTRIBUTION_MODE);
          return equalityFieldKeySelector(
              tableName, schema, equalityFields, writeParallelism, maxWriteParallelism);
        }

      default:
        throw new IllegalArgumentException(
            tableName + ": Unrecognized " + WRITE_DISTRIBUTION_MODE + ": " + mode);

View on GitHub (pinned to 86d9c8fc54)