{"record":{"id":"ad915449c14bf34a","repo":"apache/iceberg","slug":"hash-distribute-rows-by-equality-fields-even-thou-ad9154","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":"info","filePath":"flink/v2.1/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.1/flink/src/main/java/org/apache/iceberg/flink/sink/FlinkSink.java#L649-L685","documentation":"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.","triggerScenarios":"distributeDataStream encounters case RANGE with a non-empty equalityFieldIds while writing to a table with a primary key in streaming mode.","commonSituations":"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.","solutions":["Accept the fallback (hash by equality fields) - this is the safe behavior for primary-key streams","Change the table's write.distribution-mode to hash for streaming primary-key writes","If range distribution is truly needed for an append stream, clear equality fields and use the batch writer path"],"exampleFix":"// before\ntable.updateProperties().set(TableProperties.WRITE_DISTRIBUTION_MODE, \"range\").commit();\n// after\ntable.updateProperties().set(TableProperties.WRITE_DISTRIBUTION_MODE, \"hash\").commit();","handlingStrategy":"validation","validationCode":"if (\"range\".equals(table.properties().get(TableProperties.WRITE_DISTRIBUTION_MODE))\n    && !equalityFieldIds.isEmpty()) {\n  throw new IllegalStateException(\"range distribution unsupported with primary keys in Flink streaming; use hash\");\n}","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Use hash distribution mode for streaming primary-key tables","Keep batch and streaming table property profiles separate","Assert distribution mode at job startup before writing"],"tags":["flink","iceberg","distribution","primary-key"],"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"}