{"record":{"id":"9e29c0e102c62e98","repo":"apache/iceberg","slug":"hash-distribute-rows-by-equality-fields-even-thou-9e29c0","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/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.1/flink/src/main/java/org/apache/iceberg/flink/sink/IcebergSink.java#L1032-L1068","documentation":"In IcebergSink's range-distribution path, if equality fields (primary key) are set, range distribution is unsafe in a Flink streaming writer. The sink logs this warning and falls back to hash distribution keyed by equality fields for backward compatibility instead of throwing.","triggerScenarios":"IcebergSink constructor's distribution logic hits WRITE_DISTRIBUTION_MODE=range (or a sortOrder implying range) with non-empty equalityFieldIds in a streaming write.","commonSituations":"Table property write.distribution-mode=range set for batch workloads and reused by a Flink streaming upsert job; users expecting global ordering but silently receiving hash-by-PK shuffles.","solutions":["Accept the hash-by-equality-fields fallback, or switch the table property to write.distribution-mode=hash explicitly","If range distribution is required, write with a batch writer or without equality fields","Document that range mode is ignored for primary-key streaming writes"],"exampleFix":"// before\nUPDATE table SET TBLPROPERTIES ('write.distribution-mode'='range');\n// after\nUPDATE table SET TBLPROPERTIES ('write.distribution-mode'='hash');","handlingStrategy":"validation","validationCode":"if (\"range\".equals(table.properties().get(TableProperties.WRITE_DISTRIBUTION_MODE))\n    && !equalityFieldIds.isEmpty()) {\n  // streaming + PK: range is unsafe; switch to hash before building the sink\n  table.updateProperties().set(TableProperties.WRITE_DISTRIBUTION_MODE, \"hash\").commit();\n}","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Audit table properties when reusing batch tables for streaming writes","Use hash mode for all primary-key Flink streaming sinks","Log expected vs actual distribution mode at job startup"],"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"}