{"record":{"id":"1daa8dec3da8728a","repo":"apache/iceberg","slug":"hash-distribute-rows-by-equality-fields-even-thou-1daa8d","errorCode":null,"errorMessage":"Hash distribute rows by equality fields, even though {}=range is set. Range distribution for primary keys are not always safe in Flink streaming writer.","messagePattern":"Hash distribute rows by equality fields, even though (.+?)=range is set\\. Range distribution for primary keys are not always safe in Flink streaming writer\\.","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/sink/FlinkSink.java","lineNumber":667,"sourceCode":"              for (PartitionField partitionField : partitionSpec.fields()) {\n                Preconditions.checkState(\n                    equalityFieldIds.contains(partitionField.sourceId()),\n                    \"In 'hash' distribution mode with equality fields set, source column '%s' of partition field '%s' \"\n                        + \"should be included in equality fields: '%s'\",\n                    table.schema().findColumnName(partitionField.sourceId()),\n                    partitionField,\n                    equalityFieldColumns);\n              }\n              return input.keyBy(new PartitionKeySelector(partitionSpec, iSchema, flinkRowType));\n            }\n          }\n\n        case RANGE:\n          // Ideally, exception should be thrown in the combination of range distribution and\n          // equality fields. Primary key case should use hash distribution mode.\n          // Keep the current behavior of falling back to keyBy for backward compatibility.\n          if (!equalityFieldIds.isEmpty()) {\n            LOG.warn(\n                \"Hash distribute rows by equality fields, even though {}=range is set. \"\n                    + \"Range distribution for primary keys are not always safe in \"\n                    + \"Flink streaming writer.\",\n                WRITE_DISTRIBUTION_MODE);\n            return input.keyBy(\n                new EqualityFieldKeySelector(iSchema, flinkRowType, equalityFieldIds));\n          }\n\n          // range distribute by partition key or sort key if table has an SortOrder\n          Preconditions.checkState(\n              sortOrder.isSorted() || partitionSpec.isPartitioned(),\n              \"Invalid write distribution mode: range. Need to define sort order or partition spec.\");\n          if (sortOrder.isUnsorted()) {\n            sortOrder = Partitioning.sortOrderFor(partitionSpec);\n            LOG.info(\"Construct sort order from partition spec\");\n          }\n\n          LOG.info(\"Range distribute rows by sort order: {}\", sortOrder);","sourceCodeStart":649,"sourceCodeEnd":685,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/sink/FlinkSink.java#L649-L685","documentation":"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.","triggerScenarios":"distributeDataStream with DistributionMode.RANGE and non-empty equalityFieldIds — e.g. setting write.distribution-mode=range on a table with identifier fields / equalityFieldColumns configured.","commonSituations":"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.","solutions":["Accept the fallback (hash by equality fields) — this is the safe behavior.","Set write.distribution-mode=hash explicitly to make intent clear and silence the warning.","Remove equalityFieldColumns if range distribution without keys is truly intended."],"exampleFix":"// before\nprops.put(\"write.distribution-mode\", \"range\"); // table has primary keys\n// after\nprops.put(\"write.distribution-mode\", \"hash\");","handlingStrategy":"validation","validationCode":"if (mode == DistributionMode.RANGE && !eqFieldIds.isEmpty()) {\n  mode = DistributionMode.HASH; // range is unsafe with primary keys in streaming writer\n}","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Never combine range distribution with equality fields for Flink streaming writes","Use hash mode for primary-key tables; reserve range for key-less batch writes","Document distribution-mode per table in table properties"],"tags":["flink","distribution-mode","primary-key","configuration"],"backgroundTag":"conflicting-config-options","analyzedSha":"86d9c8fc543e7c56c9f624eb725f76c9baff9570","analyzedAt":"2026-09-12T00:46:39.097Z","contentChangedAt":"2026-09-12T00:46:39.097Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}