leaves12138 commented on code in PR #890:
URL: https://github.com/apache/paimon-rust/pull/890#discussion_r4058685279


##########
crates/paimon/src/table/chunk_shuffle.rs:
##########
@@ -0,0 +1,1035 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//     http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+//! Deterministic, fixed-row chunk planning for native Python reads.
+//!
+//! The shuffle intentionally matches `random.Random(seed).shuffle` in CPython.
+//! PyPaimon exposed that ordering before native planning existed, so using a
+//! different Rust RNG would silently assign different chunks to workers.
+
+use std::cmp::Ordering;
+use std::collections::{BTreeMap, HashMap};
+
+use crate::deletion_vector::{DeletionVector, DeletionVectorFactory};
+use crate::spec::{BinaryRow, DataField, DataFileMeta, Datum};
+use crate::table::source::{data_evolution_anchor_file, 
is_data_evolution_normal_file};
+use crate::table::stats_filter::group_by_overlapping_row_id;
+use crate::table::{merge_row_ranges, DataSplit, DataSplitBuilder, 
DeletionFile, RowRange, Table};
+
+/// Native chunk-shuffle configuration. The seed is stored as the unsigned
+/// little-endian 32-bit words consumed by CPython's MT19937 initializer.
+#[derive(Debug, Clone, PartialEq, Eq)]
+pub(crate) struct ChunkShuffle {
+    seed_words: Vec<u32>,
+    chunk_size: i64,
+}
+
+impl ChunkShuffle {
+    /// Build from a Python integer's decimal spelling. Negative integers use
+    /// their absolute value, matching `random.Random`.
+    pub(crate) fn from_decimal_seed(seed: &str, chunk_size: u64) -> 
crate::Result<Self> {
+        let chunk_size = i64::try_from(chunk_size).map_err(|_| 
crate::Error::DataInvalid {
+            message: format!("chunk_shuffle chunk_size {chunk_size} exceeds 
i64::MAX"),
+            source: None,
+        })?;
+        if chunk_size == 0 {
+            return Err(crate::Error::DataInvalid {
+                message: "chunk_shuffle chunk_size must be 
positive".to_string(),
+                source: None,
+            });
+        }
+        Ok(Self {
+            seed_words: decimal_seed_words(seed)?,
+            chunk_size,
+        })
+    }
+}
+
+#[derive(Debug, Clone)]
+struct InputFile {
+    file: DataFileMeta,
+    deletion_file: Option<DeletionFile>,
+}
+
+#[derive(Debug)]
+struct InputGroup {
+    partition: BinaryRow,
+    bucket: i32,
+    bucket_path: String,
+    total_buckets: i32,
+    snapshot_id: i64,
+    is_streaming: bool,
+    files: Vec<InputFile>,
+}
+
+#[derive(Debug)]
+struct AppendSegment {
+    input: InputFile,
+    ranges: Vec<RowRange>,
+}
+
+#[derive(Debug)]
+struct EvolutionSegment {
+    files: Vec<InputFile>,
+    ranges: Vec<RowRange>,
+}
+
+/// Repack planned files into shuffled, fixed-live-row chunks. Normal scan
+/// planning runs first, so partition/stats/projection pruning and 
deletion-file
+/// resolution remain centralized in `TableScan`.
+pub(crate) async fn chunk_shuffle_splits(
+    table: &Table,
+    splits: Vec<DataSplit>,
+    config: &ChunkShuffle,
+    shard: Option<(usize, usize)>,
+) -> crate::Result<Vec<DataSplit>> {
+    if !table.schema().primary_keys().is_empty() {
+        return Err(crate::Error::Unsupported {
+            message: "chunk_shuffle only supports append tables".to_string(),
+        });
+    }
+    if splits.iter().any(|split| split.row_ranges().is_some()) {
+        return Err(crate::Error::Unsupported {
+            message: "chunk_shuffle cannot combine with row-range 
selection".to_string(),
+        });
+    }
+
+    let partition_fields = partition_fields(table)?;
+    let mut groups = flatten_groups(splits)?;
+    groups.sort_by(|left, right| {
+        compare_partitions(&left.partition, &right.partition, 
&partition_fields)
+            .unwrap_or_else(|_| {
+                left.partition
+                    .to_serialized_bytes()
+                    .cmp(&right.partition.to_serialized_bytes())
+            })
+            .then_with(|| left.bucket.cmp(&right.bucket))
+    });
+
+    let data_evolution = 
table.schema().core_options().data_evolution_enabled();
+    let mut chunks = Vec::new();
+    for mut group in groups {
+        if data_evolution {
+            group.files.sort_by(|left, right| {
+                left.file
+                    .first_row_id
+                    .cmp(&right.file.first_row_id)
+                    .then_with(|| {
+                        is_data_evolution_normal_file(&right.file)
+                            .cmp(&is_data_evolution_normal_file(&left.file))
+                    })
+                    .then_with(|| 
left.file.file_name.cmp(&right.file.file_name))
+            });
+            chunks.extend(evolution_chunks(table, group, 
config.chunk_size).await?);
+        } else {
+            group
+                .files
+                .sort_by(|left, right| 
left.file.file_name.cmp(&right.file.file_name));
+            chunks.extend(append_chunks(table, group, 
config.chunk_size).await?);
+        }
+    }
+
+    PythonRandom::new(&config.seed_words).shuffle(&mut chunks);
+    if let Some((index, count)) = shard {
+        let (start, end) = shard_range(chunks.len(), index, count);
+        chunks = chunks.drain(start..end).collect();
+    }
+    Ok(chunks)
+}
+
+fn partition_fields(table: &Table) -> crate::Result<Vec<DataField>> {
+    let fields = table.schema().fields();
+    table
+        .schema()
+        .partition_keys()
+        .iter()
+        .map(|name| {
+            fields
+                .iter()
+                .find(|field| field.name() == name)
+                .cloned()
+                .ok_or_else(|| crate::Error::DataInvalid {
+                    message: format!("partition field '{name}' does not 
exist"),
+                    source: None,
+                })
+        })
+        .collect()
+}
+
+fn compare_partitions(
+    left: &BinaryRow,
+    right: &BinaryRow,
+    fields: &[DataField],
+) -> crate::Result<Ordering> {
+    for (index, field) in fields.iter().enumerate() {
+        let left = left.get_datum(index, field.data_type())?;
+        let right = right.get_datum(index, field.data_type())?;
+        // Python's key is `(value is None, value)`: non-null sorts first.
+        let ordering = match (left, right) {
+            (None, None) => Ordering::Equal,
+            (None, Some(_)) => Ordering::Greater,
+            (Some(_), None) => Ordering::Less,
+            (Some(left), Some(right)) => compare_datums(&left, &right),
+        };
+        if ordering != Ordering::Equal {
+            return Ok(ordering);
+        }
+    }
+    Ok(Ordering::Equal)
+}
+
+fn compare_datums(left: &Datum, right: &Datum) -> Ordering {
+    left.partial_cmp(right).unwrap_or_else(|| {
+        // NaN has no total order. Python's stable sort leaves incomparable
+        // values in input order; the binary representation is a deterministic
+        // fallback when manifest concurrency changed that input order.
+        left.to_string().cmp(&right.to_string())
+    })
+}
+
+fn flatten_groups(splits: Vec<DataSplit>) -> crate::Result<Vec<InputGroup>> {
+    let mut grouped: BTreeMap<(Vec<u8>, i32), InputGroup> = BTreeMap::new();
+    for split in splits {
+        let deletion_files = split.data_deletion_files();
+        let key = (split.partition().to_serialized_bytes(), split.bucket());
+        let group = grouped.entry(key).or_insert_with(|| InputGroup {
+            partition: split.partition().clone(),
+            bucket: split.bucket(),
+            bucket_path: split.bucket_path().to_string(),
+            total_buckets: split.total_buckets(),
+            snapshot_id: split.snapshot_id(),
+            is_streaming: split.is_streaming(),
+            files: Vec::new(),
+        });
+        if group.bucket_path != split.bucket_path()
+            || group.total_buckets != split.total_buckets()
+            || group.snapshot_id != split.snapshot_id()
+            || group.is_streaming != split.is_streaming()
+        {
+            return Err(crate::Error::DataInvalid {
+                message: "inconsistent split metadata within a partition 
bucket".to_string(),
+                source: None,
+            });
+        }
+        for (index, file) in split.data_files().iter().cloned().enumerate() {
+            group.files.push(InputFile {
+                file,
+                deletion_file: deletion_files
+                    .and_then(|files| files.get(index))
+                    .cloned()
+                    .flatten(),
+            });
+        }
+    }
+    Ok(grouped.into_values().collect())
+}
+
+async fn append_chunks(
+    table: &Table,
+    mut group: InputGroup,
+    chunk_size: i64,
+) -> crate::Result<Vec<DataSplit>> {
+    let row_tracking = table.schema().core_options().row_tracking_enabled();
+    let mut chunks: Vec<Vec<AppendSegment>> = Vec::new();
+    let mut current = Vec::new();
+    let mut current_rows = 0;
+
+    let inputs = std::mem::take(&mut group.files);
+    for input in inputs {
+        let mut slicer = match live_row_slicer(table, &input).await? {
+            Some(slicer) => slicer,
+            None => continue,
+        };
+        loop {
+            if current_rows == chunk_size {
+                chunks.push(std::mem::take(&mut current));
+                current_rows = 0;
+            }
+            let Some(slice) = slicer.take(chunk_size - current_rows)? else {
+                break;
+            };
+            current_rows += slice.live_rows;
+            current.push(AppendSegment {
+                input: input.clone(),
+                ranges: slice.ranges,
+            });
+        }
+    }
+    if !current.is_empty() {
+        chunks.push(current);
+    }
+
+    chunks
+        .into_iter()
+        .map(|segments| build_append_split(&group, segments, row_tracking))
+        .collect()
+}
+
+fn build_append_split(
+    group: &InputGroup,
+    segments: Vec<AppendSegment>,
+    row_tracking: bool,
+) -> crate::Result<DataSplit> {
+    let mut files = Vec::with_capacity(segments.len());
+    let mut deletion_files = Vec::with_capacity(segments.len());
+    let mut ranges = Vec::new();
+    let mut split_offset = 0;
+    for segment in segments {
+        let range_base = if row_tracking {
+            segment
+                .input
+                .file
+                .first_row_id

Review Comment:
   [P1] Align the companion Python paths with the new row-tracked range 
coordinates
   
   Using first_row_id here makes the Rust planner/reader agree, but the current 
companion #10023 still uses split-local coordinates for all non-DE append 
IndexedSplits. Its `AppendChunkShuffleSplitGenerator._chunk_to_split` builds 
ranges from split_offset, and `RawFileSplitRead.__init__` always passes ranges 
through `_split_local_row_ranges_by_file`, including on row-tracking tables.
   
   Reproduced against freshly built 7596b12 plus companion fb79b66: create a 
non-DE append table with `row-tracking.enabled=true`, commit `[0, 1, 2]` and 
`[3, 4, 5]` as two files, then use this normal public API:
   
   ```python
   builder = table.copy({
       'scan.native-plan.enabled': 'true',
       'read.native.enabled': 'false',
   }).new_read_builder()
   plan = builder.new_scan().with_chunk_shuffle(42, 3).plan()
   rows = builder.new_read().to_arrow(plan.splits())
   ```
   
   Expected IDs are `[0, 1, 2, 3, 4, 5]`; only `[0, 1, 2]` are returned. Rust 
emits `[3, 5]` for the second single-file split, but Python treats that as a 
local range for a three-row file and discards the whole chunk.
   
   The reverse supported path also fails: plan with both native options 
disabled, then feed those Python-planned splits to a reader with 
`read.native.enabled=true`. Python emits `[0, 2]` for that file, while the new 
Rust reader expects global IDs starting at 3, so it discards the same three 
rows. Pure Python and pure Rust both return all six rows.
   
   Please update the companion Python generator and raw reader (or normalize at 
the bridge) so both sides agree, and add row-tracking-without-DE coverage for 
all four planner/reader combinations. This is a current user-API result 
regression, not a request to preserve an unpublished split representation.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to