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 db380ca3 refactor(vindex): trim granule planning overhead (#892)
db380ca3 is described below
commit db380ca3d48d15f710e9474c5cb4cc11a3b3440a
Author: jerry <[email protected]>
AuthorDate: Mon Sep 21 13:40:51 2026 +0800
refactor(vindex): trim granule planning overhead (#892)
---
crates/paimon/src/arrow/format/parquet.rs | 1 -
.../table/vindex_index_build_builder/pipeline.rs | 22 +++++++---------------
2 files changed, 7 insertions(+), 16 deletions(-)
diff --git a/crates/paimon/src/arrow/format/parquet.rs
b/crates/paimon/src/arrow/format/parquet.rs
index 34cd8335..c1c8be80 100644
--- a/crates/paimon/src/arrow/format/parquet.rs
+++ b/crates/paimon/src/arrow/format/parquet.rs
@@ -98,7 +98,6 @@ impl ParquetFormatReader {
}
}
-#[derive(Debug, Clone)]
pub(crate) struct ParquetGranule {
pub(crate) first_row: i64,
pub(crate) row_count: i64,
diff --git a/crates/paimon/src/table/vindex_index_build_builder/pipeline.rs
b/crates/paimon/src/table/vindex_index_build_builder/pipeline.rs
index ebd2b8b3..26088ca3 100644
--- a/crates/paimon/src/table/vindex_index_build_builder/pipeline.rs
+++ b/crates/paimon/src/table/vindex_index_build_builder/pipeline.rs
@@ -56,14 +56,12 @@ const QUEUE_CAPACITY: usize = 2;
const BUFFER_BYTES: usize = 8 * 1024 * 1024;
const REPLAY_TARGET_BYTES: usize = 32 * 1024 * 1024;
-#[derive(Clone, Debug)]
struct Granule {
range: RowRange,
file_index: usize,
byte_ranges: Vec<Range<u64>>,
}
-#[derive(Debug)]
pub(super) struct GranulePlan {
pub(super) first: Vec<RowRange>,
pub(super) rest: Vec<RowRange>,
@@ -518,11 +516,12 @@ fn select_first(
selected.extend(unit.clone());
}
}
+ if selected.len() == granules.len() || selected.len() <
MIN_STRATA.min(granules.len()) {
+ return Ok(whole_shard());
+ }
let first_bytes = granule_bytes(granules, Some(&selected));
let total_bytes = granule_bytes(granules, None);
- if selected.len() == granules.len()
- || selected.len() < MIN_STRATA.min(granules.len())
- || total_bytes == 0
+ if total_bytes == 0
|| first_bytes.saturating_mul(FIRST_BYTES_DENOMINATOR)
> total_bytes.saturating_mul(FIRST_BYTES_NUMERATOR)
{
@@ -771,14 +770,7 @@ impl<'a> VindexIndexBuildBuilder<'a> {
granules = vec![Granule {
range: shard_range,
file_index: 0,
- byte_ranges: std::iter::once(
- 0..shard
- .files
- .iter()
- .map(|file| file.file_size.max(0) as u64)
- .sum(),
- )
- .collect(),
+ byte_ranges: Vec::new(),
}];
}
@@ -1384,7 +1376,7 @@ mod tests {
let picks = |shard: &VindexIndexShard| {
pick_indices(4_096, 512, &mut
shard_seed(shard)).collect::<Vec<_>>()
};
- assert_eq!(picks(&shard), picks(&shard.clone()));
+ assert_eq!(picks(&shard), picks(&shard));
let mut others = vec![shard.clone(); 5];
others[0].partition_bytes = vec![1, 2, 3, 4, 5, 6, 7, 8, 10];
@@ -1605,7 +1597,7 @@ mod tests {
.unwrap();
let spill = spawn_spill_writer(false).unwrap();
- for _ in 0..2 {
+ for _ in 0..QUEUE_CAPACITY {
spill
.sender
.send(SpillRecord {