This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git
The following commit(s) were added to refs/heads/main by this push:
new 157406bd feat: support append row-position selection in native scans
(#893)
157406bd is described below
commit 157406bdf6a2761dd73810d0fb1051a45291a246
Author: Jingsong Lee <[email protected]>
AuthorDate: Mon Sep 21 11:07:28 2026 +0800
feat: support append row-position selection in native scans (#893)
---
bindings/python/src/read.rs | 4 +-
bindings/python/tests/test_read.py | 94 +++++-
crates/paimon/src/table/row_position_selection.rs | 394 +++++++++++++++++++++-
crates/paimon/src/table/table_scan.rs | 51 ++-
docs/src/python-binding.md | 23 +-
5 files changed, 523 insertions(+), 43 deletions(-)
diff --git a/bindings/python/src/read.rs b/bindings/python/src/read.rs
index c6605a2b..56516bb6 100644
--- a/bindings/python/src/read.rs
+++ b/bindings/python/src/read.rs
@@ -533,7 +533,7 @@ impl PyTableScan {
#[pymethods]
impl PyTableScan {
- /// Select a half-open range of Data Evolution row positions.
+ /// Select a half-open range of append-table row positions.
fn with_row_position_slice(
mut slf: PyRefMut<'_, Self>,
start: u64,
@@ -546,7 +546,7 @@ impl PyTableScan {
Ok(slf)
}
- /// Select one Data Evolution row-position shard.
+ /// Select one balanced append-table row-position shard.
fn with_row_position_shard(
mut slf: PyRefMut<'_, Self>,
index: u64,
diff --git a/bindings/python/tests/test_read.py
b/bindings/python/tests/test_read.py
index 559b4f82..02eb9dc2 100644
--- a/bindings/python/tests/test_read.py
+++ b/bindings/python/tests/test_read.py
@@ -1201,11 +1201,99 @@ def
test_row_position_selection_validates_parameters_and_combinations():
):
with pytest.raises(ValueError, match="cannot be used
simultaneously"):
getattr(scan, method)(*args)
+
+
+def _make_append_position_table(warehouse, row_tracking=False):
+ ctx = SQLContext()
+ ctx.register_catalog("paimon", {"warehouse": warehouse})
+ ctx.sql("CREATE SCHEMA paimon.appendpos")
+ options = " WITH ('source.split.target-size' = '1b'"
+ if row_tracking:
+ options += ", 'row-tracking.enabled' = 'true'"
+ options += ")"
+ ctx.sql("CREATE TABLE paimon.appendpos.t (id INT, value STRING)" + options)
+ for start in range(0, 9, 3):
+ values = ", ".join(
+ "(%d, '%s')" % (value, chr(ord('a') + value))
+ for value in range(start, start + 3)
+ )
+ ctx.sql("INSERT INTO paimon.appendpos.t VALUES " + values)
+ return PaimonCatalog({"warehouse": warehouse}).get_table("appendpos.t")
+
+
[email protected]("row_tracking", [False, True])
+def
test_append_row_position_slices_and_shards_are_native_readable(row_tracking):
with tempfile.TemporaryDirectory() as warehouse:
- ordinary = _make_table_with_data(warehouse).new_read_builder()
+ table = _make_append_position_table(warehouse, row_tracking)
+ builder = table.new_read_builder().with_projection(["id"])
+
+ for start, end, expected in (
+ (0, 1, [0]),
+ (2, 7, [2, 3, 4, 5, 6]),
+ (7, 100, [7, 8]),
+ (20, 22, []),
+ ):
+ plan = builder.new_scan().with_row_position_slice(start,
end).plan()
+ restored = [Split.deserialize(split.serialize()) for split in
plan.splits()]
+ batches = builder.new_read().read(restored)
+ actual = pa.Table.from_batches(batches).column("id").to_pylist()
if batches else []
+ assert actual == expected
+ assert plan.snapshot_id() == 3
+
+ shards = []
+ for index, expected in enumerate(([0, 1, 2], [3, 4], [5, 6], [7, 8])):
+ plan = builder.new_scan().with_row_position_shard(index, 4).plan()
+ actual =
pa.Table.from_batches(builder.new_read().read(plan.splits()))
+ values = actual.column("id").to_pylist()
+ assert values == expected
+ shards.extend(values)
+ assert shards == list(range(9))
+
+
+def test_append_row_positions_follow_stats_pruned_file_order_before_limit():
+ with tempfile.TemporaryDirectory() as warehouse:
+ table = _make_append_position_table(warehouse)
+ builder = (
+ table.new_read_builder()
+ .with_projection(["id"])
+ .with_filter({"method": "greaterOrEqual", "field": "id",
"literals": [3]})
+ .with_limit(2)
+ )
+ plan = builder.new_scan().with_row_position_slice(1, 5).plan()
+ actual = pa.Table.from_batches(builder.new_read().read(plan.splits()))
+ # with_limit is a planning hint in the Rust API. The complete selected
+ # range must survive planning; PyPaimon enforces the final two-row
limit.
+ assert actual.column("id").to_pylist() == [4, 5, 6, 7]
+
+
+def test_append_incremental_row_positions_use_combined_delta_order():
+ with tempfile.TemporaryDirectory() as warehouse:
+ table = _make_append_position_table(warehouse)
+ builder = table.new_read_builder().with_projection(["id"])
+ for start_snapshot, end_snapshot, start, end, expected in (
+ (0, 3, 2, 7, [2, 3, 4, 5, 6]),
+ (1, 3, 1, 4, [4, 5, 6]),
+ ):
+ plan = (
+ builder.new_incremental_scan(start_snapshot, end_snapshot)
+ .with_row_position_slice(start, end)
+ .plan()
+ )
+ restored = [pickle.loads(pickle.dumps(split)) for split in
plan.splits()]
+ actual = pa.Table.from_batches(builder.new_read().read(restored))
+ assert actual.column("id").to_pylist() == expected
+
+
+def test_primary_key_row_position_selection_remains_unsupported():
+ with tempfile.TemporaryDirectory() as warehouse:
+ ctx = SQLContext()
+ ctx.register_catalog("paimon", {"warehouse": warehouse})
+ ctx.sql("CREATE SCHEMA paimon.pkpos")
+ ctx.sql("CREATE TABLE paimon.pkpos.t (id INT, value STRING, PRIMARY
KEY (id) NOT ENFORCED)")
+ table = PaimonCatalog({"warehouse": warehouse}).get_table("pkpos.t")
for method in ("with_row_position_slice", "with_row_position_shard"):
- with pytest.raises(NotImplementedError, match="data.evolution|Data
Evolution"):
- getattr(ordinary.new_scan(), method)(0, 1)
+ with pytest.raises(NotImplementedError, match="append tables"):
+ getattr(table.new_read_builder().new_scan(), method)(0, 1)
def test_incremental_row_positions_use_combined_delta_batch():
diff --git a/crates/paimon/src/table/row_position_selection.rs
b/crates/paimon/src/table/row_position_selection.rs
index 7c5b74bc..e8c62982 100644
--- a/crates/paimon/src/table/row_position_selection.rs
+++ b/crates/paimon/src/table/row_position_selection.rs
@@ -15,9 +15,9 @@
// specific language governing permissions and limitations
// under the License.
-//! Map row positions to data-evolution row IDs before group pruning.
+//! Select balanced row-position slices for data-evolution and append scans.
-use super::source::{merge_row_ranges, RowRange};
+use super::source::{merge_row_ranges, DataSplit, DataSplitBuilder, RowRange};
use crate::spec::ManifestEntry;
use crate::{Error, Result};
@@ -53,6 +53,18 @@ impl RowPositionSelection {
matches!(self, Self::Slice { .. })
}
+ fn bounds(self, total: u64) -> (u64, u64) {
+ match self {
+ Self::Slice { start, end } => (start, end.min(total)),
+ Self::Shard { index, count } => {
+ let size = total / count;
+ let remainder = total % count;
+ let start = index * size + index.min(remainder);
+ (start, start + size + u64::from(index < remainder))
+ }
+ }
+ }
+
/// Positions count the union of complete candidate file ranges. Updates
/// and blob column files share row IDs and must not multiply that count;
/// deleted rows still occupy positions until their files leave the
snapshot.
@@ -90,15 +102,7 @@ impl RowPositionSelection {
// counts so a range ending at i64::MAX cannot overflow
RowRange::count.
let count = |range: &RowRange| range.to() as u64 - range.from() as u64
+ 1;
let total: u64 = ranges.iter().map(count).sum();
- let (start, end) = match self {
- Self::Slice { start, end } => (start, end.min(total)),
- Self::Shard { index, count } => {
- let size = total / count;
- let remainder = total % count;
- let start = index * size + index.min(remainder);
- (start, start + size + u64::from(index < remainder))
- }
- };
+ let (start, end) = self.bounds(total);
if start >= end {
return Vec::new();
}
@@ -139,11 +143,188 @@ impl RowPositionSelection {
}
merge_row_ranges(result)
}
+
+ /// Select physical append rows after stats pruning and split packing.
+ ///
+ /// Append positions follow the final split/file order. Row-tracked tables
+ /// carry stable global row IDs to the reader; tables without row tracking
+ /// carry positions local to the filtered output split. Filtering files
here
+ /// avoids opening files which contain no selected rows.
+ pub(crate) fn select_append_splits(
+ self,
+ splits: Vec<DataSplit>,
+ row_tracking_enabled: bool,
+ ) -> Result<Vec<DataSplit>> {
+ let mut total = 0u64;
+ for split in &splits {
+ for file in split.data_files() {
+ let count = u64::try_from(file.row_count).map_err(|_|
Error::DataInvalid {
+ message: format!(
+ "Row-position selection requires a valid row count for
'{}'",
+ file.file_name
+ ),
+ source: None,
+ })?;
+ total = total.checked_add(count).ok_or_else(||
Error::DataInvalid {
+ message: "Row-position selection row count
overflow".to_string(),
+ source: None,
+ })?;
+ }
+ }
+ let (start, end) = self.bounds(total);
+ if start >= end {
+ return Ok(Vec::new());
+ }
+
+ let mut position = 0u64;
+ let mut selected_splits = Vec::new();
+ for split in splits {
+ let mut files = Vec::new();
+ let mut deletion_files = split.data_deletion_files().map(|_|
Vec::new());
+ let mut ranges = Vec::new();
+ let mut kept_position = 0u64;
+ let mut needs_ranges = split.row_ranges().is_some();
+
+ for (file_index, file) in split.data_files().iter().enumerate() {
+ let count = file.row_count as u64;
+ let next = position + count;
+ let selected_start = start.max(position);
+ let selected_end = end.min(next);
+ if selected_start < selected_end {
+ let from = selected_start - position;
+ let to = selected_end - position;
+ needs_ranges |= from != 0 || to != count;
+ files.push(file.clone());
+ if let (Some(source), Some(selected)) =
+ (split.data_deletion_files(), deletion_files.as_mut())
+ {
+ selected.push(source[file_index].clone());
+ }
+
+ let range = if row_tracking_enabled {
+ let first_row_id = file.first_row_id.ok_or_else(||
Error::DataInvalid {
+ message: format!(
+ "Row-position selection requires a valid first
row id for '{}'",
+ file.file_name
+ ),
+ source: None,
+ })?;
+ let from = first_row_id.checked_add(from as
i64).ok_or_else(|| {
+ Error::DataInvalid {
+ message: "Row-position selection row id
overflow".to_string(),
+ source: None,
+ }
+ })?;
+ let to = first_row_id.checked_add(to as i64 -
1).ok_or_else(|| {
+ Error::DataInvalid {
+ message: "Row-position selection row id
overflow".to_string(),
+ source: None,
+ }
+ })?;
+ RowRange::new(from, to)
+ } else {
+ let from =
+ kept_position
+ .checked_add(from)
+ .ok_or_else(|| Error::DataInvalid {
+ message: "Row-position selection split
offset overflow"
+ .to_string(),
+ source: None,
+ })?;
+ let to = kept_position.checked_add(to -
1).ok_or_else(|| {
+ Error::DataInvalid {
+ message: "Row-position selection split offset
overflow".to_string(),
+ source: None,
+ }
+ })?;
+ RowRange::new(
+ i64::try_from(from).map_err(|_| Error::DataInvalid
{
+ message: "Row-position selection split offset
exceeds i64"
+ .to_string(),
+ source: None,
+ })?,
+ i64::try_from(to).map_err(|_| Error::DataInvalid {
+ message: "Row-position selection split offset
exceeds i64"
+ .to_string(),
+ source: None,
+ })?,
+ )
+ };
+ ranges.push(range);
+ kept_position += count;
+ }
+ position = next;
+ }
+
+ if files.is_empty() {
+ if position >= end {
+ break;
+ }
+ continue;
+ }
+
+ let ranges = if let Some(existing) = split.row_ranges() {
+ intersect_ranges(&ranges, existing)
+ } else {
+ merge_row_ranges(ranges)
+ };
+ if ranges.is_empty() {
+ if position >= end {
+ break;
+ }
+ continue;
+ }
+
+ let mut builder = DataSplitBuilder::new()
+ .with_snapshot(split.snapshot_id())
+ .with_partition(split.partition().clone())
+ .with_bucket(split.bucket())
+ .with_bucket_path(split.bucket_path().to_string())
+ .with_total_buckets(split.total_buckets())
+ .with_data_files(files)
+ .with_raw_convertible(split.raw_convertible())
+ .with_streaming(split.is_streaming());
+ if let Some(deletion_files) = deletion_files {
+ builder = builder.with_data_deletion_files(deletion_files);
+ }
+ if needs_ranges {
+ builder = builder.with_row_ranges(ranges);
+ }
+ selected_splits.push(builder.build()?);
+ if position >= end {
+ break;
+ }
+ }
+ Ok(selected_splits)
+ }
+}
+
+fn intersect_ranges(left: &[RowRange], right: &[RowRange]) -> Vec<RowRange> {
+ let left = merge_row_ranges(left.to_vec());
+ let right = merge_row_ranges(right.to_vec());
+ let mut result = Vec::new();
+ let (mut left_index, mut right_index) = (0, 0);
+ while left_index < left.len() && right_index < right.len() {
+ if let Some(range) =
+ left[left_index].intersect_inclusive(right[right_index].from(),
right[right_index].to())
+ {
+ result.push(range);
+ }
+ if left[left_index].to() <= right[right_index].to() {
+ left_index += 1;
+ } else {
+ right_index += 1;
+ }
+ }
+ merge_row_ranges(result)
}
#[cfg(test)]
mod tests {
use super::*;
+ use crate::spec::stats::BinaryTableStats;
+ use crate::spec::{BinaryRow, DataFileMeta};
+ use crate::table::source::DeletionFile;
fn ranges(pairs: &[(i64, i64)]) -> Vec<RowRange> {
pairs
@@ -152,6 +333,197 @@ mod tests {
.collect()
}
+ fn file(name: &str, row_count: i64, first_row_id: Option<i64>) ->
DataFileMeta {
+ DataFileMeta {
+ file_name: name.to_string(),
+ file_size: 100,
+ row_count,
+ min_key: Vec::new(),
+ max_key: Vec::new(),
+ key_stats: BinaryTableStats::new(Vec::new(), Vec::new(),
Vec::new()),
+ value_stats: BinaryTableStats::new(Vec::new(), Vec::new(),
Vec::new()),
+ min_sequence_number: 0,
+ max_sequence_number: 0,
+ schema_id: 0,
+ level: 0,
+ extra_files: Vec::new(),
+ creation_time: None,
+ delete_row_count: None,
+ embedded_index: None,
+ first_row_id,
+ write_cols: None,
+ external_path: None,
+ file_source: None,
+ value_stats_cols: None,
+ column_max_sequence_numbers: None,
+ }
+ }
+
+ fn split(
+ snapshot: i64,
+ files: Vec<DataFileMeta>,
+ deletion_files: Option<Vec<Option<DeletionFile>>>,
+ row_ranges: Option<Vec<RowRange>>,
+ ) -> DataSplit {
+ let mut builder = DataSplitBuilder::new()
+ .with_snapshot(snapshot)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(2)
+ .with_bucket_path("memory:/append/bucket-2".to_string())
+ .with_total_buckets(4)
+ .with_data_files(files)
+ .with_raw_convertible(true)
+ .with_streaming(true);
+ if let Some(deletion_files) = deletion_files {
+ builder = builder.with_data_deletion_files(deletion_files);
+ }
+ if let Some(row_ranges) = row_ranges {
+ builder = builder.with_row_ranges(row_ranges);
+ }
+ builder.build().unwrap()
+ }
+
+ fn pairs(ranges: Option<&[RowRange]>) -> Option<Vec<(i64, i64)>> {
+ ranges.map(|ranges| {
+ ranges
+ .iter()
+ .map(|range| (range.from(), range.to()))
+ .collect()
+ })
+ }
+
+ #[test]
+ fn append_slice_filters_files_and_rebases_local_ranges() {
+ let deletion = DeletionFile::new("b.dv".to_string(), 0, 10, Some(1));
+ let input = split(
+ 7,
+ vec![file("a.parquet", 3, None), file("b.parquet", 4, None)],
+ Some(vec![None, Some(deletion.clone())]),
+ None,
+ );
+
+ let selected = RowPositionSelection::slice(4, 6)
+ .unwrap()
+ .select_append_splits(vec![input], false)
+ .unwrap();
+
+ assert_eq!(selected.len(), 1);
+ let selected = &selected[0];
+ assert_eq!(
+ selected
+ .data_files()
+ .iter()
+ .map(|file| file.file_name.as_str())
+ .collect::<Vec<_>>(),
+ vec!["b.parquet"]
+ );
+ assert_eq!(pairs(selected.row_ranges()), Some(vec![(1, 2)]));
+ assert_eq!(
+ selected.data_deletion_files(),
+ Some([Some(deletion)].as_slice())
+ );
+ assert_eq!(selected.snapshot_id(), 7);
+ assert_eq!(selected.bucket(), 2);
+ assert_eq!(selected.total_buckets(), 4);
+ assert!(selected.is_streaming());
+ assert!(selected.raw_convertible());
+ }
+
+ #[test]
+ fn append_slice_uses_global_ids_when_row_tracking_is_enabled() {
+ let input = split(
+ 1,
+ vec![
+ file("a.parquet", 3, Some(100)),
+ file("b.parquet", 4, Some(200)),
+ ],
+ None,
+ None,
+ );
+
+ let selected = RowPositionSelection::slice(2, 5)
+ .unwrap()
+ .select_append_splits(vec![input], true)
+ .unwrap();
+ assert_eq!(selected.len(), 1);
+ assert_eq!(
+ pairs(selected[0].row_ranges()),
+ Some(vec![(102, 102), (200, 201)])
+ );
+
+ let restored =
+
DataSplit::deserialize_split_v1(&selected[0].serialize_split_v1().unwrap()).unwrap();
+ assert_eq!(
+ pairs(restored.row_ranges()),
+ pairs(selected[0].row_ranges())
+ );
+ }
+
+ #[test]
+ fn append_shards_cover_physical_rows_once_across_splits() {
+ let inputs = vec![
+ split(
+ 1,
+ vec![file("a.parquet", 3, None), file("b.parquet", 4, None)],
+ None,
+ None,
+ ),
+ split(1, vec![file("c.parquet", 2, None)], None, None),
+ ];
+ let mut counts = Vec::new();
+ for index in 0..4 {
+ let selected = RowPositionSelection::shard(index, 4)
+ .unwrap()
+ .select_append_splits(inputs.clone(), false)
+ .unwrap();
+
counts.push(selected.iter().map(DataSplit::row_count).sum::<i64>());
+ }
+ assert_eq!(counts, vec![3, 2, 2, 2]);
+ assert_eq!(counts.iter().sum::<i64>(), 9);
+ }
+
+ #[test]
+ fn append_selection_intersects_existing_global_ranges() {
+ let input = split(
+ 1,
+ vec![
+ file("a.parquet", 3, Some(100)),
+ file("b.parquet", 4, Some(200)),
+ ],
+ None,
+ Some(ranges(&[(201, 203)])),
+ );
+ let selected = RowPositionSelection::slice(2, 5)
+ .unwrap()
+ .select_append_splits(vec![input], true)
+ .unwrap();
+ assert_eq!(pairs(selected[0].row_ranges()), Some(vec![(201, 201)]));
+ }
+
+ #[test]
+ fn append_selection_rejects_unknown_counts_and_missing_row_ids() {
+ let unknown = split(
+ 1,
+ vec![file(
+ "unknown.parquet",
+ DataFileMeta::ROW_COUNT_UNKNOWN,
+ None,
+ )],
+ None,
+ None,
+ );
+ assert!(RowPositionSelection::slice(0, 1)
+ .unwrap()
+ .select_append_splits(vec![unknown], false)
+ .is_err());
+
+ let missing_id = split(1, vec![file("missing.parquet", 1, None)],
None, None);
+ assert!(RowPositionSelection::slice(0, 1)
+ .unwrap()
+ .select_append_splits(vec![missing_id], true)
+ .is_err());
+ }
+
#[test]
fn slice_counts_shared_row_ids_once_and_skips_gaps() {
let candidates = ranges(&[(10, 12), (0, 3), (1, 2), (11, 12), (20,
21)]);
diff --git a/crates/paimon/src/table/table_scan.rs
b/crates/paimon/src/table/table_scan.rs
index 46864b8e..3ccb3562 100644
--- a/crates/paimon/src/table/table_scan.rs
+++ b/crates/paimon/src/table/table_scan.rs
@@ -1126,13 +1126,16 @@ impl<'a> TableScan<'a> {
}
}
- /// Select a half-open range of logical positions in a data-evolution
snapshot.
- /// Positions are assigned before group statistics, projection and
deletion vectors.
+ /// Select a half-open range of logical positions in an append snapshot.
+ ///
+ /// Data-evolution positions are assigned before group statistics,
projection
+ /// and deletion vectors. Ordinary append positions follow the final
+ /// stats-pruned split and file order.
pub fn with_row_position_slice(self, start: u64, end: u64) ->
crate::Result<Self> {
self.with_row_position_selection(RowPositionSelection::slice(start,
end)?)
}
- /// Select a balanced contiguous shard of data-evolution row positions.
+ /// Select a balanced contiguous shard of append row positions.
pub fn with_row_position_shard(self, index: u64, count: u64) ->
crate::Result<Self> {
self.with_row_position_selection(RowPositionSelection::shard(index,
count)?)
}
@@ -1235,9 +1238,7 @@ impl<'a> TableScan<'a> {
fn with_row_position_selection(self, selection: RowPositionSelection) ->
crate::Result<Self> {
match self.0 {
- TableScanKind::Paimon(mut scan)
- if scan.table.schema().core_options().data_evolution_enabled()
=>
- {
+ TableScanKind::Paimon(mut scan) if
scan.table.schema().primary_keys().is_empty() => {
if scan.chunk_shuffle().is_some() {
return Err(crate::Error::DataInvalid {
message:
@@ -1271,7 +1272,7 @@ impl<'a> TableScan<'a> {
Ok(Self(TableScanKind::Paimon(scan)))
}
_ => Err(crate::Error::Unsupported {
- message: "row-position selection requires a data-evolution
table".into(),
+ message: "row-position selection only supports append
tables".into(),
}),
}
}
@@ -2338,14 +2339,20 @@ impl<'a> PaimonTableScan<'a> {
let open_file_cost = core_options.source_split_open_file_cost();
let partition_keys = self.table.schema().partition_keys();
- // Assign row positions using the full candidate file ranges, before
- // group stats, projection or DVs change visible rows. Intersect
explicit
- // and index-selected ranges only after assigning the positional range.
- let effective_row_ranges = if let Some(selection) =
self.row_position_selection() {
- Some(selection.select(&entries, effective_row_ranges.as_deref())?)
- } else {
- effective_row_ranges
- };
+ let row_position_selection = self.row_position_selection();
+ let append_row_position_selection = (!data_evolution_enabled)
+ .then_some(row_position_selection)
+ .flatten();
+ // Data-evolution positions use the union of full candidate row-id
+ // ranges, before group stats, projection or DVs change visible rows.
+ // Ordinary append positions are selected from completed splits below,
+ // after stats pruning and packing establish their physical order.
+ let effective_row_ranges =
+ if let Some(selection) = row_position_selection.filter(|_|
data_evolution_enabled) {
+ Some(selection.select(&entries,
effective_row_ranges.as_deref())?)
+ } else {
+ effective_row_ranges
+ };
if effective_row_ranges.as_ref().is_some_and(Vec::is_empty) {
if let Some(trace) = trace {
trace.record_final_plan(0, 0, 0);
@@ -2463,7 +2470,10 @@ impl<'a> PaimonTableScan<'a> {
.index_file_in_data_file_dir();
let mut data_file_field_ids_cache = DataFileFieldIdsCache::new();
- let can_push_down_limit =
self.can_push_down_limit_hint(effective_row_ranges.as_deref());
+ // Positional distribution precedes LIMIT. Building too few ordinary
+ // append splits here could starve a later slice or shard.
+ let can_push_down_limit = append_row_position_selection.is_none()
+ && self.can_push_down_limit_hint(effective_row_ranges.as_deref());
let mut limit_accumulator = match self.limit {
Some(limit) if limit > 0 && can_push_down_limit => {
Some(LimitPushdownAccumulator::new(limit))
@@ -2721,6 +2731,11 @@ impl<'a> PaimonTableScan<'a> {
let split_candidates_built = splits.len();
(splits, split_candidates_built, false)
};
+ let splits = if let Some(selection) = append_row_position_selection {
+ selection.select_append_splits(splits,
core_options.row_tracking_enabled())?
+ } else {
+ splits
+ };
let splits = if let Some(config) = self.chunk_shuffle() {
chunk_shuffle_splits(self.table, splits, config,
self.shard()).await?
} else {
@@ -3819,13 +3834,13 @@ mod tests {
}
#[test]
- fn test_row_position_selection_rejects_unsupported_and_mixed_modes() {
+ fn test_row_position_selection_accepts_append_and_rejects_mixed_modes() {
let append = limit_test_table();
assert!(append
.new_read_builder()
.new_scan()
.with_row_position_shard(0, 1)
- .is_err());
+ .is_ok());
let table = data_evolution_test_table(
"memory:/row_position_selection_validation",
two_column_schema(0, "id", "name"),
diff --git a/docs/src/python-binding.md b/docs/src/python-binding.md
index cf61c7ac..0c6d287e 100644
--- a/docs/src/python-binding.md
+++ b/docs/src/python-binding.md
@@ -169,22 +169,27 @@ to historical events. Builder filters, projections, and
limits still apply.
both batch and streaming splits to Java binary encoding, preserving the
streaming
flag.
-For Data Evolution tables, select half-open row positions or one balanced shard
-on a scan:
+For append and Data Evolution tables, select half-open row positions or one
+balanced shard on a scan:
```python
plan = rb.new_scan().with_row_position_slice(10, 20).plan()
plan = rb.new_scan().with_row_position_shard(1, 4).plan()
```
-Positions count candidate rows before explicit/global-index range pruning,
-group statistics, projection, and deletion-vector filtering. Column updates
+For ordinary append tables, positions follow the final stats-pruned split and
+file order. Row-tracked tables encode selected stable row IDs; tables without
+row tracking encode positions local to each filtered output split. Unselected
+files are removed, so the reader does not open them. For Data Evolution,
+positions count candidate rows before explicit/global-index range pruning,
+group statistics, projection, and deletion-vector filtering; column updates
sharing row IDs count once. Explicit `with_row_ranges` and index-selected
ranges
-intersect the positions after assignment, even when they exclude earlier
files. Slices require `start < end`; shards require a positive
-count and `0 <= index < count`. Slice and shard selection are mutually
exclusive
-and may also be applied to `new_incremental_scan` results, where positions
count
-the combined APPEND-delta batch. The selection is encoded in the returned
splits
-and survives serialization.
+intersect the positions after assignment, even when they exclude earlier files.
+Slices require `start < end`; shards require a positive count and
+`0 <= index < count`. Slice and shard selection are mutually exclusive and may
+also be applied to `new_incremental_scan` results, where positions count the
+combined APPEND-delta batch. The selection is encoded in the returned splits
+and survives serialization and direct native reads.
Deletion-vector reads accept both Java and Python Avro array-item schemas.
Historical Python bucket-local references are resolved from the table's index