{"record":{"id":"4c878c1e1d276eae","repo":"apache/iceberg","slug":"hash-distribute-rows-by-equality-fields-even-thou-4c878c","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/IcebergSink.java","lineNumber":1050,"sourceCode":"    return Optional.ofNullable(flinkWriteConf.writeParallelism()).orElseGet(input::getParallelism);\n  }\n\n  private DataStream<RowData> distributeDataStreamByRangeDistributionMode(\n      DataStream<RowData> input,\n      Schema iSchema,\n      PartitionSpec partitionSpec,\n      SortOrder sortOrderParam) {\n\n    int writerParallelism = resolveWriterParallelism(input);\n\n    // needed because of checkStyle not allowing us to change the value of an argument\n    SortOrder sortOrder = sortOrderParam;\n\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(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);\n    StatisticsOrRecordTypeInformation statisticsOrRecordTypeInformation =","sourceCodeStart":1032,"sourceCodeEnd":1068,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergSink.java#L1032-L1068","documentation":"In IcebergSink, requesting range distribution while equality fields (primary keys) are configured triggers this warning: the sink ignores the range request and hashes rows by the equality fields instead, because range distribution with primary keys is not always safe in the Flink streaming writer. Kept as a fallback rather than an error for backward compatibility.","triggerScenarios":"IcebergSink.builder().distributionMode(DistributionMode.RANGE) (or write.distribution-mode=range) with non-empty equalityFieldIds, in the range-distribution path of the builder.","commonSituations":"Batch tuning configs (range) reused for streaming jobs on tables with primary keys; expectation that range mode balances write parallelism while the table's identifier fields force hash-by-key.","solutions":["Set distributionMode to HASH explicitly so intent matches actual behavior.","Remove equalityFieldColumns if range distribution without keys is intended.","Leave as-is if the hash-by-key fallback is acceptable; it is safe and only logs."],"exampleFix":"// before\nIcebergSink.builder().distributionMode(DistributionMode.RANGE)... // equality fields set\n// after\nIcebergSink.builder().distributionMode(DistributionMode.HASH)...","handlingStrategy":"validation","validationCode":"if (mode == DistributionMode.RANGE && !equalityFieldIds.isEmpty()) {\n  mode = DistributionMode.HASH; // match actual sink behavior\n}","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Use hash mode whenever equality fields/primary keys are configured","Keep batch range-mode configs separate from streaming sink configs","Silence ambiguity by setting distribution mode explicitly at the builder"],"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"}