From 04d5dc5b2e58478a1ccba7550bfcc3802ae9c07a Mon Sep 17 00:00:00 2001 From: wangzheyan Date: Mon, 10 Aug 2026 23:08:28 +0800 Subject: [PATCH 1/6] feat(java): expose aggregate fragment summary --- java/lance-jni/src/blocking_dataset.rs | 38 +++++++++++- java/src/main/java/org/lance/Dataset.java | 17 ++++++ .../main/java/org/lance/FragmentSummary.java | 61 +++++++++++++++++++ .../src/test/java/org/lance/FragmentTest.java | 13 ++++ rust/lance/src/dataset/statistics.rs | 51 ++++++++++++++++ 5 files changed, 179 insertions(+), 1 deletion(-) create mode 100644 java/src/main/java/org/lance/FragmentSummary.java diff --git a/java/lance-jni/src/blocking_dataset.rs b/java/lance-jni/src/blocking_dataset.rs index 9d93afc95ba..13f373b20bb 100644 --- a/java/lance-jni/src/blocking_dataset.rs +++ b/java/lance-jni/src/blocking_dataset.rs @@ -35,7 +35,7 @@ use lance::dataset::cleanup::{ }; use lance::dataset::optimize::{CompactionOptions as RustCompactionOptions, compact_files}; use lance::dataset::refs::{Ref, TagContents}; -use lance::dataset::statistics::{DataStatistics, DatasetStatisticsExt}; +use lance::dataset::statistics::{DataStatistics, DatasetStatisticsExt, FragmentSummary}; use lance::dataset::transaction::{Operation, Transaction}; use lance::dataset::{ ColumnAlteration, CommitBuilder, Dataset, NewColumnTransform, ProjectionRequest, ReadParams, @@ -785,6 +785,22 @@ impl IntoJava for Version { } } +impl IntoJava for FragmentSummary { + fn into_java<'a>(self, env: &mut JNIEnv<'a>) -> Result> { + Ok(env.new_object( + "org/lance/FragmentSummary", + "(JJJJJ)V", + &[ + JValue::Long(self.fragment_count as i64), + JValue::Long(self.min_rows_per_fragment as i64), + JValue::Long(self.max_rows_per_fragment as i64), + JValue::Long(self.min_data_files_per_fragment as i64), + JValue::Long(self.max_data_files_per_fragment as i64), + ], + )?) + } +} + fn attach_native_dataset<'local>( env: &mut JNIEnv<'local>, dataset: BlockingDataset, @@ -1531,6 +1547,26 @@ pub extern "system" fn Java_org_lance_Dataset_nativeGetFragmentStatistics<'a>( ok_or_throw!(env, inner_get_fragment_statistics(&mut env, jdataset)) } +#[unsafe(no_mangle)] +pub extern "system" fn Java_org_lance_Dataset_nativeGetFragmentSummary<'a>( + mut env: JNIEnv<'a>, + jdataset: JObject, +) -> JObject<'a> { + ok_or_throw!(env, inner_get_fragment_summary(&mut env, jdataset)) +} + +fn inner_get_fragment_summary<'local>( + env: &mut JNIEnv<'local>, + jdataset: JObject, +) -> Result> { + let summary = { + let dataset = + unsafe { env.get_rust_field::<_, _, BlockingDataset>(jdataset, NATIVE_DATASET) }?; + dataset.inner.fragment_summary() + }; + summary.into_java(env) +} + /// Returns per-fragment statistics flattened as [id0, rowCount0, dataFileNum0, id1, ...]. /// /// Row count semantics match Java `FragmentMetadata.getNumRows()`: diff --git a/java/src/main/java/org/lance/Dataset.java b/java/src/main/java/org/lance/Dataset.java index 02b7c81ce39..4f73c39f4f4 100644 --- a/java/src/main/java/org/lance/Dataset.java +++ b/java/src/main/java/org/lance/Dataset.java @@ -1379,6 +1379,23 @@ public FragmentStatistics getFragmentStatistics() { private native long[] nativeGetFragmentStatistics(); + /** + * Get aggregate statistics for all fragments in this dataset version. + * + *

The aggregation runs in native code and returns a fixed-size result, avoiding per-fragment + * Java objects and arrays. + * + * @return aggregate fragment statistics + */ + public FragmentSummary getFragmentSummary() { + try (LockManager.ReadLock readLock = lockManager.acquireReadLock()) { + Preconditions.checkArgument(nativeDatasetHandle != 0, "Dataset is closed"); + return nativeGetFragmentSummary(); + } + } + + private native FragmentSummary nativeGetFragmentSummary(); + /** * Gets the arrow schema of the dataset. * diff --git a/java/src/main/java/org/lance/FragmentSummary.java b/java/src/main/java/org/lance/FragmentSummary.java new file mode 100644 index 00000000000..c755f2ef424 --- /dev/null +++ b/java/src/main/java/org/lance/FragmentSummary.java @@ -0,0 +1,61 @@ +/* + * Licensed 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. + */ +package org.lance; + +/** Aggregate statistics for the fragments in a dataset version. */ +public final class FragmentSummary { + private final long fragmentCount; + private final long minRowsPerFragment; + private final long maxRowsPerFragment; + private final long minDataFilesPerFragment; + private final long maxDataFilesPerFragment; + + FragmentSummary( + long fragmentCount, + long minRowsPerFragment, + long maxRowsPerFragment, + long minDataFilesPerFragment, + long maxDataFilesPerFragment) { + this.fragmentCount = fragmentCount; + this.minRowsPerFragment = minRowsPerFragment; + this.maxRowsPerFragment = maxRowsPerFragment; + this.minDataFilesPerFragment = minDataFilesPerFragment; + this.maxDataFilesPerFragment = maxDataFilesPerFragment; + } + + /** Number of fragments. */ + public long getFragmentCount() { + return fragmentCount; + } + + /** Minimum number of live rows in a fragment, or 0 when there are no fragments. */ + public long getMinRowsPerFragment() { + return minRowsPerFragment; + } + + /** Maximum number of live rows in a fragment, or 0 when there are no fragments. */ + public long getMaxRowsPerFragment() { + return maxRowsPerFragment; + } + + /** Minimum number of data files in a fragment, or 0 when there are no fragments. */ + public long getMinDataFilesPerFragment() { + return minDataFilesPerFragment; + } + + /** Maximum number of data files in a fragment, or 0 when there are no fragments. */ + public long getMaxDataFilesPerFragment() { + return maxDataFilesPerFragment; + } +} diff --git a/java/src/test/java/org/lance/FragmentTest.java b/java/src/test/java/org/lance/FragmentTest.java index 879475cdab0..a2d37c32642 100644 --- a/java/src/test/java/org/lance/FragmentTest.java +++ b/java/src/test/java/org/lance/FragmentTest.java @@ -441,6 +441,13 @@ void testFragmentStatistics(@TempDir Path tempDir) { stats.getDataFileNums()); assertEquals(30, Arrays.stream(stats.getRowCounts()).sum()); + + FragmentSummary summary = dataset.getFragmentSummary(); + assertEquals(2, summary.getFragmentCount()); + assertEquals(9, summary.getMinRowsPerFragment()); + assertEquals(21, summary.getMaxRowsPerFragment()); + assertEquals(1, summary.getMinDataFilesPerFragment()); + assertEquals(1, summary.getMaxDataFilesPerFragment()); } } } @@ -453,6 +460,12 @@ void testFragmentStatisticsOnEmptyDataset(@TempDir Path tempDir) { new TestUtils.SimpleTestDataset(allocator, datasetPath); try (Dataset dataset = testDataset.createEmptyDataset()) { assertEquals(0, dataset.getFragmentStatistics().size()); + FragmentSummary summary = dataset.getFragmentSummary(); + assertEquals(0, summary.getFragmentCount()); + assertEquals(0, summary.getMinRowsPerFragment()); + assertEquals(0, summary.getMaxRowsPerFragment()); + assertEquals(0, summary.getMinDataFilesPerFragment()); + assertEquals(0, summary.getMaxDataFilesPerFragment()); } } } diff --git a/rust/lance/src/dataset/statistics.rs b/rust/lance/src/dataset/statistics.rs index 47b224e1c49..fecb9b4cfa8 100644 --- a/rust/lance/src/dataset/statistics.rs +++ b/rust/lance/src/dataset/statistics.rs @@ -17,6 +17,57 @@ use super::overlay::{collect_overlay_stale_frags, overlaid_fragments}; use super::{Dataset, fragment::FileFragment}; use crate::index::{DatasetIndexExt, DatasetIndexInternalExt}; +/// Aggregate statistics for the fragments in a dataset version. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct FragmentSummary { + /// Number of fragments. + pub fragment_count: u64, + /// Minimum number of live rows in a fragment, or 0 when the dataset has no fragments. + pub min_rows_per_fragment: u64, + /// Maximum number of live rows in a fragment, or 0 when the dataset has no fragments. + pub max_rows_per_fragment: u64, + /// Minimum number of data files in a fragment, or 0 when the dataset has no fragments. + pub min_data_files_per_fragment: u64, + /// Maximum number of data files in a fragment, or 0 when the dataset has no fragments. + pub max_data_files_per_fragment: u64, +} + +impl Dataset { + /// Aggregate fragment statistics from the loaded manifest in one pass. + pub fn fragment_summary(&self) -> FragmentSummary { + let mut summary = FragmentSummary { + fragment_count: self.fragments().len() as u64, + min_rows_per_fragment: u64::MAX, + max_rows_per_fragment: 0, + min_data_files_per_fragment: u64::MAX, + max_data_files_per_fragment: 0, + }; + + for fragment in self.fragments().iter() { + let live_rows = fragment.physical_rows.unwrap_or(0).saturating_sub( + fragment + .deletion_file + .as_ref() + .and_then(|deletion_file| deletion_file.num_deleted_rows) + .unwrap_or(0), + ) as u64; + let data_file_count = fragment.files.len() as u64; + summary.min_rows_per_fragment = summary.min_rows_per_fragment.min(live_rows); + summary.max_rows_per_fragment = summary.max_rows_per_fragment.max(live_rows); + summary.min_data_files_per_fragment = + summary.min_data_files_per_fragment.min(data_file_count); + summary.max_data_files_per_fragment = + summary.max_data_files_per_fragment.max(data_file_count); + } + + if summary.fragment_count == 0 { + summary.min_rows_per_fragment = 0; + summary.min_data_files_per_fragment = 0; + } + summary + } +} + /// Statistics about a single field in the dataset pub struct FieldStatistics { /// Id of the field From a92d342b7d6340c2da39c78b64c9e324d60f0311 Mon Sep 17 00:00:00 2001 From: wangzheyan Date: Tue, 11 Aug 2026 01:10:46 +0800 Subject: [PATCH 2/6] fix(java): preserve unknown fragment row counts --- java/lance-jni/src/blocking_dataset.rs | 7 +- .../main/java/org/lance/FragmentSummary.java | 36 ++++- .../src/test/java/org/lance/FragmentTest.java | 28 +++- rust/lance/src/dataset/statistics.rs | 145 ++++++++++++++---- 4 files changed, 173 insertions(+), 43 deletions(-) diff --git a/java/lance-jni/src/blocking_dataset.rs b/java/lance-jni/src/blocking_dataset.rs index 13f373b20bb..629069aee5f 100644 --- a/java/lance-jni/src/blocking_dataset.rs +++ b/java/lance-jni/src/blocking_dataset.rs @@ -789,11 +789,12 @@ impl IntoJava for FragmentSummary { fn into_java<'a>(self, env: &mut JNIEnv<'a>) -> Result> { Ok(env.new_object( "org/lance/FragmentSummary", - "(JJJJJ)V", + "(JJJJJJ)V", &[ JValue::Long(self.fragment_count as i64), - JValue::Long(self.min_rows_per_fragment as i64), - JValue::Long(self.max_rows_per_fragment as i64), + JValue::Long(self.min_rows_per_fragment.unwrap_or(0) as i64), + JValue::Long(self.max_rows_per_fragment.unwrap_or(0) as i64), + JValue::Long(self.unknown_row_count_fragment_count as i64), JValue::Long(self.min_data_files_per_fragment as i64), JValue::Long(self.max_data_files_per_fragment as i64), ], diff --git a/java/src/main/java/org/lance/FragmentSummary.java b/java/src/main/java/org/lance/FragmentSummary.java index c755f2ef424..814daf527d6 100644 --- a/java/src/main/java/org/lance/FragmentSummary.java +++ b/java/src/main/java/org/lance/FragmentSummary.java @@ -13,11 +13,14 @@ */ package org.lance; +import java.util.OptionalLong; + /** Aggregate statistics for the fragments in a dataset version. */ public final class FragmentSummary { private final long fragmentCount; private final long minRowsPerFragment; private final long maxRowsPerFragment; + private final long unknownRowCountFragmentCount; private final long minDataFilesPerFragment; private final long maxDataFilesPerFragment; @@ -25,11 +28,13 @@ public final class FragmentSummary { long fragmentCount, long minRowsPerFragment, long maxRowsPerFragment, + long unknownRowCountFragmentCount, long minDataFilesPerFragment, long maxDataFilesPerFragment) { this.fragmentCount = fragmentCount; this.minRowsPerFragment = minRowsPerFragment; this.maxRowsPerFragment = maxRowsPerFragment; + this.unknownRowCountFragmentCount = unknownRowCountFragmentCount; this.minDataFilesPerFragment = minDataFilesPerFragment; this.maxDataFilesPerFragment = maxDataFilesPerFragment; } @@ -39,14 +44,33 @@ public long getFragmentCount() { return fragmentCount; } - /** Minimum number of live rows in a fragment, or 0 when there are no fragments. */ - public long getMinRowsPerFragment() { - return minRowsPerFragment; + /** + * Minimum number of live rows in a fragment. + * + * @return empty when any fragment has an unknown live-row count; otherwise the minimum, or 0 when + * there are no fragments + */ + public OptionalLong getMinRowsPerFragment() { + return unknownRowCountFragmentCount == 0 + ? OptionalLong.of(minRowsPerFragment) + : OptionalLong.empty(); + } + + /** + * Maximum number of live rows in a fragment. + * + * @return empty when any fragment has an unknown live-row count; otherwise the maximum, or 0 when + * there are no fragments + */ + public OptionalLong getMaxRowsPerFragment() { + return unknownRowCountFragmentCount == 0 + ? OptionalLong.of(maxRowsPerFragment) + : OptionalLong.empty(); } - /** Maximum number of live rows in a fragment, or 0 when there are no fragments. */ - public long getMaxRowsPerFragment() { - return maxRowsPerFragment; + /** Number of fragments whose live row count is unknown. */ + public long getUnknownRowCountFragmentCount() { + return unknownRowCountFragmentCount; } /** Minimum number of data files in a fragment, or 0 when there are no fragments. */ diff --git a/java/src/test/java/org/lance/FragmentTest.java b/java/src/test/java/org/lance/FragmentTest.java index a2d37c32642..8b00b39989f 100644 --- a/java/src/test/java/org/lance/FragmentTest.java +++ b/java/src/test/java/org/lance/FragmentTest.java @@ -444,8 +444,9 @@ void testFragmentStatistics(@TempDir Path tempDir) { FragmentSummary summary = dataset.getFragmentSummary(); assertEquals(2, summary.getFragmentCount()); - assertEquals(9, summary.getMinRowsPerFragment()); - assertEquals(21, summary.getMaxRowsPerFragment()); + assertEquals(9, summary.getMinRowsPerFragment().orElseThrow(AssertionError::new)); + assertEquals(21, summary.getMaxRowsPerFragment().orElseThrow(AssertionError::new)); + assertEquals(0, summary.getUnknownRowCountFragmentCount()); assertEquals(1, summary.getMinDataFilesPerFragment()); assertEquals(1, summary.getMaxDataFilesPerFragment()); } @@ -462,11 +463,30 @@ void testFragmentStatisticsOnEmptyDataset(@TempDir Path tempDir) { assertEquals(0, dataset.getFragmentStatistics().size()); FragmentSummary summary = dataset.getFragmentSummary(); assertEquals(0, summary.getFragmentCount()); - assertEquals(0, summary.getMinRowsPerFragment()); - assertEquals(0, summary.getMaxRowsPerFragment()); + assertEquals(0, summary.getMinRowsPerFragment().orElseThrow(AssertionError::new)); + assertEquals(0, summary.getMaxRowsPerFragment().orElseThrow(AssertionError::new)); + assertEquals(0, summary.getUnknownRowCountFragmentCount()); assertEquals(0, summary.getMinDataFilesPerFragment()); assertEquals(0, summary.getMaxDataFilesPerFragment()); } } } + + @Test + void testFragmentSummaryPreservesUnknownRowsFromHistoricalManifest() { + String historicalPath = + Path.of("..", "test_data", "v0.7.5", "with_deletions") + .toAbsolutePath() + .normalize() + .toString(); + try (RootAllocator allocator = new RootAllocator(Long.MAX_VALUE); + Dataset dataset = Dataset.open(historicalPath, allocator)) { + FragmentSummary summary = dataset.getFragmentSummary(); + assertTrue(summary.getUnknownRowCountFragmentCount() > 0); + assertTrue(summary.getMinRowsPerFragment().isEmpty()); + assertTrue(summary.getMaxRowsPerFragment().isEmpty()); + assertTrue(summary.getMinDataFilesPerFragment() > 0); + assertTrue(summary.getMaxDataFilesPerFragment() > 0); + } + } } diff --git a/rust/lance/src/dataset/statistics.rs b/rust/lance/src/dataset/statistics.rs index fecb9b4cfa8..aa917c5e7c3 100644 --- a/rust/lance/src/dataset/statistics.rs +++ b/rust/lance/src/dataset/statistics.rs @@ -22,10 +22,12 @@ use crate::index::{DatasetIndexExt, DatasetIndexInternalExt}; pub struct FragmentSummary { /// Number of fragments. pub fragment_count: u64, - /// Minimum number of live rows in a fragment, or 0 when the dataset has no fragments. - pub min_rows_per_fragment: u64, - /// Maximum number of live rows in a fragment, or 0 when the dataset has no fragments. - pub max_rows_per_fragment: u64, + /// Minimum number of live rows in a fragment, or `None` when any row count is unknown. + pub min_rows_per_fragment: Option, + /// Maximum number of live rows in a fragment, or `None` when any row count is unknown. + pub max_rows_per_fragment: Option, + /// Number of fragments whose live row count is unknown. + pub unknown_row_count_fragment_count: u64, /// Minimum number of data files in a fragment, or 0 when the dataset has no fragments. pub min_data_files_per_fragment: u64, /// Maximum number of data files in a fragment, or 0 when the dataset has no fragments. @@ -35,36 +37,119 @@ pub struct FragmentSummary { impl Dataset { /// Aggregate fragment statistics from the loaded manifest in one pass. pub fn fragment_summary(&self) -> FragmentSummary { - let mut summary = FragmentSummary { - fragment_count: self.fragments().len() as u64, - min_rows_per_fragment: u64::MAX, - max_rows_per_fragment: 0, - min_data_files_per_fragment: u64::MAX, - max_data_files_per_fragment: 0, - }; + summarize_fragments(self.fragments()) + } +} - for fragment in self.fragments().iter() { - let live_rows = fragment.physical_rows.unwrap_or(0).saturating_sub( - fragment - .deletion_file - .as_ref() - .and_then(|deletion_file| deletion_file.num_deleted_rows) - .unwrap_or(0), - ) as u64; - let data_file_count = fragment.files.len() as u64; - summary.min_rows_per_fragment = summary.min_rows_per_fragment.min(live_rows); - summary.max_rows_per_fragment = summary.max_rows_per_fragment.max(live_rows); - summary.min_data_files_per_fragment = - summary.min_data_files_per_fragment.min(data_file_count); - summary.max_data_files_per_fragment = - summary.max_data_files_per_fragment.max(data_file_count); +fn summarize_fragments(fragments: &[lance_table::format::Fragment]) -> FragmentSummary { + let mut min_rows_per_fragment = u64::MAX; + let mut max_rows_per_fragment = 0; + let mut min_data_files_per_fragment = u64::MAX; + let mut max_data_files_per_fragment = 0; + let mut unknown_row_count_fragment_count = 0; + + for fragment in fragments { + match fragment.num_rows() { + Some(live_rows) => { + min_rows_per_fragment = min_rows_per_fragment.min(live_rows as u64); + max_rows_per_fragment = max_rows_per_fragment.max(live_rows as u64); + } + None => unknown_row_count_fragment_count += 1, } + let data_file_count = fragment.files.len() as u64; + min_data_files_per_fragment = min_data_files_per_fragment.min(data_file_count); + max_data_files_per_fragment = max_data_files_per_fragment.max(data_file_count); + } + + if fragments.is_empty() { + min_rows_per_fragment = 0; + min_data_files_per_fragment = 0; + } + + let row_counts_complete = unknown_row_count_fragment_count == 0; + FragmentSummary { + fragment_count: fragments.len() as u64, + min_rows_per_fragment: row_counts_complete.then_some(min_rows_per_fragment), + max_rows_per_fragment: row_counts_complete.then_some(max_rows_per_fragment), + unknown_row_count_fragment_count, + min_data_files_per_fragment, + max_data_files_per_fragment, + } +} - if summary.fragment_count == 0 { - summary.min_rows_per_fragment = 0; - summary.min_data_files_per_fragment = 0; +#[cfg(test)] +mod fragment_summary_tests { + use lance_file::version::ConcreteFileVersion; + use lance_table::format::{DeletionFile, DeletionFileType, Fragment}; + + use super::summarize_fragments; + + fn fragment(id: u64, rows: Option, file_count: usize) -> Fragment { + let mut fragment = Fragment::new(id); + fragment.physical_rows = rows; + for file_idx in 0..file_count { + fragment.add_file( + format!("{id}-{file_idx}.lance"), + vec![file_idx as i32], + vec![file_idx as i32], + ConcreteFileVersion::V2_1, + None, + ); } - summary + fragment + } + + #[test] + fn test_fragment_summary_preserves_unknown_row_counts() { + let known = fragment(0, Some(10), 1); + let unknown_physical_rows = fragment(1, None, 2); + let mut unknown_deletions = fragment(2, Some(30), 3); + unknown_deletions.deletion_file = Some(DeletionFile { + read_version: 1, + id: 1, + file_type: DeletionFileType::Array, + num_deleted_rows: None, + base_id: None, + }); + + let summary = summarize_fragments(&[known, unknown_physical_rows, unknown_deletions]); + assert_eq!(summary.fragment_count, 3); + assert_eq!(summary.min_rows_per_fragment, None); + assert_eq!(summary.max_rows_per_fragment, None); + assert_eq!(summary.unknown_row_count_fragment_count, 2); + assert_eq!(summary.min_data_files_per_fragment, 1); + assert_eq!(summary.max_data_files_per_fragment, 3); + } + + #[test] + fn test_fragment_summary_uses_live_rows() { + let mut deleted = fragment(0, Some(10), 1); + deleted.deletion_file = Some(DeletionFile { + read_version: 1, + id: 1, + file_type: DeletionFileType::Array, + num_deleted_rows: Some(4), + base_id: None, + }); + let other = fragment(1, Some(20), 2); + + let summary = summarize_fragments(&[deleted, other]); + assert_eq!(summary.min_rows_per_fragment, Some(6)); + assert_eq!(summary.max_rows_per_fragment, Some(20)); + assert_eq!(summary.unknown_row_count_fragment_count, 0); + assert_eq!(summary.min_data_files_per_fragment, 1); + assert_eq!(summary.max_data_files_per_fragment, 2); + } + + #[test] + fn test_fragment_summary_empty() { + let summary = summarize_fragments(&[]); + assert_eq!(summary.fragment_count, 0); + assert_eq!(summary.min_rows_per_fragment, Some(0)); + assert_eq!(summary.max_rows_per_fragment, Some(0)); + assert_eq!(summary.unknown_row_count_fragment_count, 0); + assert_eq!(summary.min_data_files_per_fragment, 0); + assert_eq!(summary.max_data_files_per_fragment, 0); } } From 226857c17a37803a3c10e91027a17d359bf0d400 Mon Sep 17 00:00:00 2001 From: wangzheyan Date: Tue, 11 Aug 2026 02:01:06 +0800 Subject: [PATCH 3/6] fix(java): reject fragment summaries with unknown rows --- java/lance-jni/src/blocking_dataset.rs | 9 +-- java/src/main/java/org/lance/Dataset.java | 4 + .../main/java/org/lance/FragmentSummary.java | 36 ++------- .../src/test/java/org/lance/FragmentTest.java | 21 ++---- rust/lance/src/dataset/statistics.rs | 75 +++++++++---------- 5 files changed, 59 insertions(+), 86 deletions(-) diff --git a/java/lance-jni/src/blocking_dataset.rs b/java/lance-jni/src/blocking_dataset.rs index 629069aee5f..c66f8572fc4 100644 --- a/java/lance-jni/src/blocking_dataset.rs +++ b/java/lance-jni/src/blocking_dataset.rs @@ -789,12 +789,11 @@ impl IntoJava for FragmentSummary { fn into_java<'a>(self, env: &mut JNIEnv<'a>) -> Result> { Ok(env.new_object( "org/lance/FragmentSummary", - "(JJJJJJ)V", + "(JJJJJ)V", &[ JValue::Long(self.fragment_count as i64), - JValue::Long(self.min_rows_per_fragment.unwrap_or(0) as i64), - JValue::Long(self.max_rows_per_fragment.unwrap_or(0) as i64), - JValue::Long(self.unknown_row_count_fragment_count as i64), + JValue::Long(self.min_rows_per_fragment as i64), + JValue::Long(self.max_rows_per_fragment as i64), JValue::Long(self.min_data_files_per_fragment as i64), JValue::Long(self.max_data_files_per_fragment as i64), ], @@ -1563,7 +1562,7 @@ fn inner_get_fragment_summary<'local>( let summary = { let dataset = unsafe { env.get_rust_field::<_, _, BlockingDataset>(jdataset, NATIVE_DATASET) }?; - dataset.inner.fragment_summary() + dataset.inner.fragment_summary()? }; summary.into_java(env) } diff --git a/java/src/main/java/org/lance/Dataset.java b/java/src/main/java/org/lance/Dataset.java index 4f73c39f4f4..9cb3ea7620d 100644 --- a/java/src/main/java/org/lance/Dataset.java +++ b/java/src/main/java/org/lance/Dataset.java @@ -1385,7 +1385,11 @@ public FragmentStatistics getFragmentStatistics() { *

The aggregation runs in native code and returns a fixed-size result, avoiding per-fragment * Java objects and arrays. * + *

Every fragment must contain enough metadata to determine its live row count. Some legacy + * datasets do not contain this metadata and are not supported by this method. + * * @return aggregate fragment statistics + * @throws RuntimeException if a fragment is missing required row-count metadata */ public FragmentSummary getFragmentSummary() { try (LockManager.ReadLock readLock = lockManager.acquireReadLock()) { diff --git a/java/src/main/java/org/lance/FragmentSummary.java b/java/src/main/java/org/lance/FragmentSummary.java index 814daf527d6..c755f2ef424 100644 --- a/java/src/main/java/org/lance/FragmentSummary.java +++ b/java/src/main/java/org/lance/FragmentSummary.java @@ -13,14 +13,11 @@ */ package org.lance; -import java.util.OptionalLong; - /** Aggregate statistics for the fragments in a dataset version. */ public final class FragmentSummary { private final long fragmentCount; private final long minRowsPerFragment; private final long maxRowsPerFragment; - private final long unknownRowCountFragmentCount; private final long minDataFilesPerFragment; private final long maxDataFilesPerFragment; @@ -28,13 +25,11 @@ public final class FragmentSummary { long fragmentCount, long minRowsPerFragment, long maxRowsPerFragment, - long unknownRowCountFragmentCount, long minDataFilesPerFragment, long maxDataFilesPerFragment) { this.fragmentCount = fragmentCount; this.minRowsPerFragment = minRowsPerFragment; this.maxRowsPerFragment = maxRowsPerFragment; - this.unknownRowCountFragmentCount = unknownRowCountFragmentCount; this.minDataFilesPerFragment = minDataFilesPerFragment; this.maxDataFilesPerFragment = maxDataFilesPerFragment; } @@ -44,33 +39,14 @@ public long getFragmentCount() { return fragmentCount; } - /** - * Minimum number of live rows in a fragment. - * - * @return empty when any fragment has an unknown live-row count; otherwise the minimum, or 0 when - * there are no fragments - */ - public OptionalLong getMinRowsPerFragment() { - return unknownRowCountFragmentCount == 0 - ? OptionalLong.of(minRowsPerFragment) - : OptionalLong.empty(); - } - - /** - * Maximum number of live rows in a fragment. - * - * @return empty when any fragment has an unknown live-row count; otherwise the maximum, or 0 when - * there are no fragments - */ - public OptionalLong getMaxRowsPerFragment() { - return unknownRowCountFragmentCount == 0 - ? OptionalLong.of(maxRowsPerFragment) - : OptionalLong.empty(); + /** Minimum number of live rows in a fragment, or 0 when there are no fragments. */ + public long getMinRowsPerFragment() { + return minRowsPerFragment; } - /** Number of fragments whose live row count is unknown. */ - public long getUnknownRowCountFragmentCount() { - return unknownRowCountFragmentCount; + /** Maximum number of live rows in a fragment, or 0 when there are no fragments. */ + public long getMaxRowsPerFragment() { + return maxRowsPerFragment; } /** Minimum number of data files in a fragment, or 0 when there are no fragments. */ diff --git a/java/src/test/java/org/lance/FragmentTest.java b/java/src/test/java/org/lance/FragmentTest.java index 8b00b39989f..cecab4e6197 100644 --- a/java/src/test/java/org/lance/FragmentTest.java +++ b/java/src/test/java/org/lance/FragmentTest.java @@ -444,9 +444,8 @@ void testFragmentStatistics(@TempDir Path tempDir) { FragmentSummary summary = dataset.getFragmentSummary(); assertEquals(2, summary.getFragmentCount()); - assertEquals(9, summary.getMinRowsPerFragment().orElseThrow(AssertionError::new)); - assertEquals(21, summary.getMaxRowsPerFragment().orElseThrow(AssertionError::new)); - assertEquals(0, summary.getUnknownRowCountFragmentCount()); + assertEquals(9, summary.getMinRowsPerFragment()); + assertEquals(21, summary.getMaxRowsPerFragment()); assertEquals(1, summary.getMinDataFilesPerFragment()); assertEquals(1, summary.getMaxDataFilesPerFragment()); } @@ -463,9 +462,8 @@ void testFragmentStatisticsOnEmptyDataset(@TempDir Path tempDir) { assertEquals(0, dataset.getFragmentStatistics().size()); FragmentSummary summary = dataset.getFragmentSummary(); assertEquals(0, summary.getFragmentCount()); - assertEquals(0, summary.getMinRowsPerFragment().orElseThrow(AssertionError::new)); - assertEquals(0, summary.getMaxRowsPerFragment().orElseThrow(AssertionError::new)); - assertEquals(0, summary.getUnknownRowCountFragmentCount()); + assertEquals(0, summary.getMinRowsPerFragment()); + assertEquals(0, summary.getMaxRowsPerFragment()); assertEquals(0, summary.getMinDataFilesPerFragment()); assertEquals(0, summary.getMaxDataFilesPerFragment()); } @@ -473,7 +471,7 @@ void testFragmentStatisticsOnEmptyDataset(@TempDir Path tempDir) { } @Test - void testFragmentSummaryPreservesUnknownRowsFromHistoricalManifest() { + void testFragmentSummaryRejectsUnknownRowsFromHistoricalManifest() { String historicalPath = Path.of("..", "test_data", "v0.7.5", "with_deletions") .toAbsolutePath() @@ -481,12 +479,9 @@ void testFragmentSummaryPreservesUnknownRowsFromHistoricalManifest() { .toString(); try (RootAllocator allocator = new RootAllocator(Long.MAX_VALUE); Dataset dataset = Dataset.open(historicalPath, allocator)) { - FragmentSummary summary = dataset.getFragmentSummary(); - assertTrue(summary.getUnknownRowCountFragmentCount() > 0); - assertTrue(summary.getMinRowsPerFragment().isEmpty()); - assertTrue(summary.getMaxRowsPerFragment().isEmpty()); - assertTrue(summary.getMinDataFilesPerFragment() > 0); - assertTrue(summary.getMaxDataFilesPerFragment() > 0); + RuntimeException error = assertThrows(RuntimeException.class, dataset::getFragmentSummary); + assertTrue(error.getMessage().contains("Fragment summary requires")); + assertTrue(error.getMessage().contains("fragment")); } } } diff --git a/rust/lance/src/dataset/statistics.rs b/rust/lance/src/dataset/statistics.rs index aa917c5e7c3..db318566ae8 100644 --- a/rust/lance/src/dataset/statistics.rs +++ b/rust/lance/src/dataset/statistics.rs @@ -22,12 +22,10 @@ use crate::index::{DatasetIndexExt, DatasetIndexInternalExt}; pub struct FragmentSummary { /// Number of fragments. pub fragment_count: u64, - /// Minimum number of live rows in a fragment, or `None` when any row count is unknown. - pub min_rows_per_fragment: Option, - /// Maximum number of live rows in a fragment, or `None` when any row count is unknown. - pub max_rows_per_fragment: Option, - /// Number of fragments whose live row count is unknown. - pub unknown_row_count_fragment_count: u64, + /// Minimum number of live rows in a fragment, or 0 when the dataset has no fragments. + pub min_rows_per_fragment: u64, + /// Maximum number of live rows in a fragment, or 0 when the dataset has no fragments. + pub max_rows_per_fragment: u64, /// Minimum number of data files in a fragment, or 0 when the dataset has no fragments. pub min_data_files_per_fragment: u64, /// Maximum number of data files in a fragment, or 0 when the dataset has no fragments. @@ -36,26 +34,29 @@ pub struct FragmentSummary { impl Dataset { /// Aggregate fragment statistics from the loaded manifest in one pass. - pub fn fragment_summary(&self) -> FragmentSummary { + /// + /// Returns an error for legacy fragments that do not contain enough metadata to determine + /// their live row count. + pub fn fragment_summary(&self) -> Result { summarize_fragments(self.fragments()) } } -fn summarize_fragments(fragments: &[lance_table::format::Fragment]) -> FragmentSummary { +fn summarize_fragments(fragments: &[lance_table::format::Fragment]) -> Result { let mut min_rows_per_fragment = u64::MAX; let mut max_rows_per_fragment = 0; let mut min_data_files_per_fragment = u64::MAX; let mut max_data_files_per_fragment = 0; - let mut unknown_row_count_fragment_count = 0; for fragment in fragments { - match fragment.num_rows() { - Some(live_rows) => { - min_rows_per_fragment = min_rows_per_fragment.min(live_rows as u64); - max_rows_per_fragment = max_rows_per_fragment.max(live_rows as u64); - } - None => unknown_row_count_fragment_count += 1, - } + let live_rows = fragment.num_rows().ok_or_else(|| { + Error::internal(format!( + "Fragment summary requires physical row count and deletion count in fragment metadata, but fragment {} is missing required row-count metadata. Rewrite the dataset with a current Lance version to populate it", + fragment.id + )) + })? as u64; + min_rows_per_fragment = min_rows_per_fragment.min(live_rows); + max_rows_per_fragment = max_rows_per_fragment.max(live_rows); let data_file_count = fragment.files.len() as u64; min_data_files_per_fragment = min_data_files_per_fragment.min(data_file_count); max_data_files_per_fragment = max_data_files_per_fragment.max(data_file_count); @@ -66,19 +67,18 @@ fn summarize_fragments(fragments: &[lance_table::format::Fragment]) -> FragmentS min_data_files_per_fragment = 0; } - let row_counts_complete = unknown_row_count_fragment_count == 0; - FragmentSummary { + Ok(FragmentSummary { fragment_count: fragments.len() as u64, - min_rows_per_fragment: row_counts_complete.then_some(min_rows_per_fragment), - max_rows_per_fragment: row_counts_complete.then_some(max_rows_per_fragment), - unknown_row_count_fragment_count, + min_rows_per_fragment, + max_rows_per_fragment, min_data_files_per_fragment, max_data_files_per_fragment, - } + }) } #[cfg(test)] mod fragment_summary_tests { + use lance_core::Error; use lance_file::version::ConcreteFileVersion; use lance_table::format::{DeletionFile, DeletionFileType, Fragment}; @@ -100,7 +100,7 @@ mod fragment_summary_tests { } #[test] - fn test_fragment_summary_preserves_unknown_row_counts() { + fn test_fragment_summary_rejects_unknown_row_counts() { let known = fragment(0, Some(10), 1); let unknown_physical_rows = fragment(1, None, 2); let mut unknown_deletions = fragment(2, Some(30), 3); @@ -112,13 +112,14 @@ mod fragment_summary_tests { base_id: None, }); - let summary = summarize_fragments(&[known, unknown_physical_rows, unknown_deletions]); - assert_eq!(summary.fragment_count, 3); - assert_eq!(summary.min_rows_per_fragment, None); - assert_eq!(summary.max_rows_per_fragment, None); - assert_eq!(summary.unknown_row_count_fragment_count, 2); - assert_eq!(summary.min_data_files_per_fragment, 1); - assert_eq!(summary.max_data_files_per_fragment, 3); + let physical_rows_error = + summarize_fragments(&[known.clone(), unknown_physical_rows]).unwrap_err(); + assert!(matches!(&physical_rows_error, Error::Internal { .. })); + assert!(physical_rows_error.to_string().contains("fragment 1")); + + let deletion_count_error = summarize_fragments(&[known, unknown_deletions]).unwrap_err(); + assert!(matches!(&deletion_count_error, Error::Internal { .. })); + assert!(deletion_count_error.to_string().contains("fragment 2")); } #[test] @@ -133,21 +134,19 @@ mod fragment_summary_tests { }); let other = fragment(1, Some(20), 2); - let summary = summarize_fragments(&[deleted, other]); - assert_eq!(summary.min_rows_per_fragment, Some(6)); - assert_eq!(summary.max_rows_per_fragment, Some(20)); - assert_eq!(summary.unknown_row_count_fragment_count, 0); + let summary = summarize_fragments(&[deleted, other]).unwrap(); + assert_eq!(summary.min_rows_per_fragment, 6); + assert_eq!(summary.max_rows_per_fragment, 20); assert_eq!(summary.min_data_files_per_fragment, 1); assert_eq!(summary.max_data_files_per_fragment, 2); } #[test] fn test_fragment_summary_empty() { - let summary = summarize_fragments(&[]); + let summary = summarize_fragments(&[]).unwrap(); assert_eq!(summary.fragment_count, 0); - assert_eq!(summary.min_rows_per_fragment, Some(0)); - assert_eq!(summary.max_rows_per_fragment, Some(0)); - assert_eq!(summary.unknown_row_count_fragment_count, 0); + assert_eq!(summary.min_rows_per_fragment, 0); + assert_eq!(summary.max_rows_per_fragment, 0); assert_eq!(summary.min_data_files_per_fragment, 0); assert_eq!(summary.max_data_files_per_fragment, 0); } From c9b5aba771d9d3ce31e7e41553dcf7c216b33282 Mon Sep 17 00:00:00 2001 From: wangzheyan Date: Tue, 11 Aug 2026 11:26:55 +0800 Subject: [PATCH 4/6] bench(java): compare high-scale fragment aggregation --- rust/lance/Cargo.toml | 4 + rust/lance/benches/fragment_summary.rs | 117 +++++++++++++++++++++++++ 2 files changed, 121 insertions(+) create mode 100644 rust/lance/benches/fragment_summary.rs diff --git a/rust/lance/Cargo.toml b/rust/lance/Cargo.toml index 64253478000..180bb82363c 100644 --- a/rust/lance/Cargo.toml +++ b/rust/lance/Cargo.toml @@ -207,6 +207,10 @@ harness = false name = "count_pushdown" harness = false +[[bench]] +name = "fragment_summary" +harness = false + [[bench]] name = "vector_index" harness = false diff --git a/rust/lance/benches/fragment_summary.rs b/rust/lance/benches/fragment_summary.rs new file mode 100644 index 00000000000..8fa45d99389 --- /dev/null +++ b/rust/lance/benches/fragment_summary.rs @@ -0,0 +1,117 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +//! Compares fixed-size fragment aggregation with materializing per-fragment JNI payload data. +//! +//! The flattened baseline mirrors the native work behind Java +//! `Dataset.getFragmentStatistics()`: it allocates three `i64` values per fragment before JNI +//! copies and Java-side array splitting. `Dataset.fragment_summary()` returns five scalars +//! regardless of fragment count. +//! +//! At 100,000 fragments, the flattened path allocates a 2.4 MB native vector, followed by a +//! 2.4 MB Java `long[]` JNI copy and 1.6 MB across the three final Java primitive arrays. The +//! aggregate path creates one fixed-size Java object. +//! +//! ```text +//! cargo bench -p lance --bench fragment_summary +//! ``` + +use std::hint::black_box; + +use arrow_schema::{DataType, Field as ArrowField, Schema as ArrowSchema}; +use criterion::{Criterion, Throughput, criterion_group, criterion_main}; +use lance::Dataset; +use lance::dataset::transaction::Operation; +use lance_core::utils::tempfile::TempStrDir; +use lance_table::format::Fragment; + +const NUM_FRAGMENTS: usize = 100_000; + +struct Fixture { + _data_dir: TempStrDir, + dataset: Dataset, +} + +impl Fixture { + async fn open() -> Self { + let data_dir = TempStrDir::default(); + let schema = + lance_core::datatypes::Schema::try_from(&ArrowSchema::new(vec![ArrowField::new( + "value", + DataType::Int64, + false, + )])) + .unwrap(); + let fragments = (0..NUM_FRAGMENTS) + .map(|id| { + let mut fragment = Fragment::new(id as u64); + fragment.physical_rows = Some(1_000 + id % 100); + fragment + }) + .collect(); + let operation = Operation::Overwrite { + fragments, + schema, + config_upsert_values: None, + initial_bases: None, + }; + let dataset = Dataset::commit( + data_dir.as_str(), + operation, + None, + None, + None, + Default::default(), + false, + ) + .await + .unwrap(); + + Self { + _data_dir: data_dir, + dataset, + } + } +} + +fn flattened_fragment_statistics(dataset: &Dataset) -> Vec { + let fragments = dataset.fragments(); + let mut statistics = Vec::with_capacity(fragments.len() * 3); + for fragment in fragments.iter() { + let physical_rows = fragment.physical_rows.unwrap_or(0) as i64; + let deleted_rows = fragment + .deletion_file + .as_ref() + .and_then(|deletion_file| deletion_file.num_deleted_rows) + .unwrap_or(0) as i64; + statistics.push(fragment.id as i64); + statistics.push(physical_rows - deleted_rows); + statistics.push(fragment.files.len() as i64); + } + statistics +} + +fn bench_fragment_statistics(c: &mut Criterion) { + let runtime = tokio::runtime::Runtime::new().unwrap(); + let fixture = runtime.block_on(Fixture::open()); + + let summary = fixture.dataset.fragment_summary().unwrap(); + assert_eq!(summary.fragment_count, NUM_FRAGMENTS as u64); + assert_eq!( + flattened_fragment_statistics(&fixture.dataset).len(), + NUM_FRAGMENTS * 3 + ); + + let mut group = c.benchmark_group("fragment_statistics/100k_fragments"); + group.throughput(Throughput::Elements(NUM_FRAGMENTS as u64)); + group.bench_function("aggregate_summary", |b| { + b.iter(|| black_box(fixture.dataset.fragment_summary().unwrap())) + }); + group.bench_function("materialize_flattened_statistics", |b| { + b.iter(|| black_box(flattened_fragment_statistics(&fixture.dataset))) + }); + group.finish(); +} + +criterion_group!(benches, bench_fragment_statistics); +criterion_main!(benches); From 7ec30d5a5df7003f848b86650ab425cdfc738368 Mon Sep 17 00:00:00 2001 From: wangzheyan Date: Thu, 13 Aug 2026 20:21:33 +0800 Subject: [PATCH 5/6] perf(java): reduce fragment statistics allocations --- java/lance-jni/src/blocking_dataset.rs | 106 ++++++------ java/src/main/java/org/lance/Dataset.java | 42 +---- .../main/java/org/lance/FragmentSummary.java | 61 ------- .../src/test/java/org/lance/FragmentTest.java | 59 +++++-- rust/lance/Cargo.toml | 2 +- rust/lance/benches/fragment_statistics.rs | 152 ++++++++++++++++++ rust/lance/benches/fragment_summary.rs | 117 -------------- rust/lance/src/dataset/statistics.rs | 135 ---------------- 8 files changed, 253 insertions(+), 421 deletions(-) delete mode 100644 java/src/main/java/org/lance/FragmentSummary.java create mode 100644 rust/lance/benches/fragment_statistics.rs delete mode 100644 rust/lance/benches/fragment_summary.rs diff --git a/java/lance-jni/src/blocking_dataset.rs b/java/lance-jni/src/blocking_dataset.rs index f12c2363134..618e6255eef 100644 --- a/java/lance-jni/src/blocking_dataset.rs +++ b/java/lance-jni/src/blocking_dataset.rs @@ -35,7 +35,7 @@ use lance::dataset::cleanup::{ }; use lance::dataset::optimize::{CompactionOptions as RustCompactionOptions, compact_files}; use lance::dataset::refs::{Ref, TagContents}; -use lance::dataset::statistics::{DataStatistics, DatasetStatisticsExt, FragmentSummary}; +use lance::dataset::statistics::{DataStatistics, DatasetStatisticsExt}; use lance::dataset::transaction::{Operation, Transaction}; use lance::dataset::{ ColumnAlteration, CommitBuilder, Dataset, NewColumnTransform, ProjectionRequest, ReadParams, @@ -785,22 +785,6 @@ impl IntoJava for Version { } } -impl IntoJava for FragmentSummary { - fn into_java<'a>(self, env: &mut JNIEnv<'a>) -> Result> { - Ok(env.new_object( - "org/lance/FragmentSummary", - "(JJJJJ)V", - &[ - JValue::Long(self.fragment_count as i64), - JValue::Long(self.min_rows_per_fragment as i64), - JValue::Long(self.max_rows_per_fragment as i64), - JValue::Long(self.min_data_files_per_fragment as i64), - JValue::Long(self.max_data_files_per_fragment as i64), - ], - )?) - } -} - fn attach_native_dataset<'local>( env: &mut JNIEnv<'local>, dataset: BlockingDataset, @@ -1547,27 +1531,7 @@ pub extern "system" fn Java_org_lance_Dataset_nativeGetFragmentStatistics<'a>( ok_or_throw!(env, inner_get_fragment_statistics(&mut env, jdataset)) } -#[unsafe(no_mangle)] -pub extern "system" fn Java_org_lance_Dataset_nativeGetFragmentSummary<'a>( - mut env: JNIEnv<'a>, - jdataset: JObject, -) -> JObject<'a> { - ok_or_throw!(env, inner_get_fragment_summary(&mut env, jdataset)) -} - -fn inner_get_fragment_summary<'local>( - env: &mut JNIEnv<'local>, - jdataset: JObject, -) -> Result> { - let summary = { - let dataset = - unsafe { env.get_rust_field::<_, _, BlockingDataset>(jdataset, NATIVE_DATASET) }?; - dataset.inner.fragment_summary()? - }; - summary.into_java(env) -} - -/// Returns per-fragment statistics flattened as [id0, rowCount0, dataFileNum0, id1, ...]. +/// Returns per-fragment statistics in their final Java primitive arrays. /// /// Row count semantics match Java `FragmentMetadata.getNumRows()`: /// physical rows minus deleted rows, with absent values treated as 0. @@ -1576,28 +1540,62 @@ fn inner_get_fragment_statistics<'local>( env: &mut JNIEnv<'local>, jdataset: JObject, ) -> Result> { - let stats: Vec = { + // Three 4096-entry typed buffers use 64 KiB while keeping JNI calls amortized. + const CHUNK_SIZE: usize = 4096; + + let fragments = { let dataset = unsafe { env.get_rust_field::<_, _, BlockingDataset>(jdataset, NATIVE_DATASET) }?; - let fragments = dataset.inner.get_fragments(); - let mut stats = Vec::with_capacity(fragments.len() * 3); - for f in fragments.iter() { - let meta = f.metadata(); - let physical_rows = meta.physical_rows.unwrap_or(0) as i64; - let deleted_rows = meta + dataset.inner.fragments().clone() + }; + let fragment_count = i32::try_from(fragments.len()).map_err(|_| { + Error::runtime_error(format!( + "Fragment statistics contain {} fragments, exceeding the Java array limit of {}", + fragments.len(), + i32::MAX + )) + })?; + let ids = env.new_int_array(fragment_count)?; + let row_counts = env.new_long_array(fragment_count)?; + let data_file_nums = env.new_int_array(fragment_count)?; + + let chunk_capacity = fragments.len().min(CHUNK_SIZE); + let mut id_chunk = Vec::with_capacity(chunk_capacity); + let mut row_count_chunk = Vec::with_capacity(chunk_capacity); + let mut data_file_num_chunk = Vec::with_capacity(chunk_capacity); + + for (chunk_index, fragment_chunk) in fragments.chunks(CHUNK_SIZE).enumerate() { + id_chunk.clear(); + row_count_chunk.clear(); + data_file_num_chunk.clear(); + + for fragment in fragment_chunk { + let physical_rows = fragment.physical_rows.unwrap_or(0) as i64; + let deleted_rows = fragment .deletion_file .as_ref() - .and_then(|d| d.num_deleted_rows) + .and_then(|deletion_file| deletion_file.num_deleted_rows) .unwrap_or(0) as i64; - stats.push(f.id() as i64); - stats.push(physical_rows - deleted_rows); - stats.push(meta.files.len() as i64); + id_chunk.push(fragment.id as i32); + row_count_chunk.push(physical_rows - deleted_rows); + data_file_num_chunk.push(fragment.files.len() as i32); } - stats - }; - let jarray = env.new_long_array(stats.len() as i32)?; - env.set_long_array_region(&jarray, 0, &stats)?; - Ok(jarray.into()) + + let offset = (chunk_index * CHUNK_SIZE) as i32; + env.set_int_array_region(&ids, offset, &id_chunk)?; + env.set_long_array_region(&row_counts, offset, &row_count_chunk)?; + env.set_int_array_region(&data_file_nums, offset, &data_file_num_chunk)?; + } + + Ok(env.new_object( + "org/lance/FragmentStatistics", + "([I[J[I)V", + &[ + JValue::Object(&ids), + JValue::Object(&row_counts), + JValue::Object(&data_file_nums), + ], + )?) } #[unsafe(no_mangle)] diff --git a/java/src/main/java/org/lance/Dataset.java b/java/src/main/java/org/lance/Dataset.java index 9cb3ea7620d..4cae2aafc63 100644 --- a/java/src/main/java/org/lance/Dataset.java +++ b/java/src/main/java/org/lance/Dataset.java @@ -1353,52 +1353,20 @@ public List getFragments() { * Get per-fragment statistics for all fragments in this dataset version. * *

Unlike {@link #getFragments()}, this is a metadata-only bulk operation: no per-fragment Java - * objects are materialized, making it suitable for planning over datasets with a very large - * number of fragments. Row counts match {@link FragmentMetadata#getNumRows()} (physical rows - * minus deleted rows). + * objects are materialized, and native code fills the returned primitive arrays directly. This + * makes it suitable for planning over datasets with a very large number of fragments. Row counts + * match {@link FragmentMetadata#getNumRows()} (physical rows minus deleted rows). * * @return per-fragment statistics as parallel arrays, in manifest order */ public FragmentStatistics getFragmentStatistics() { try (LockManager.ReadLock readLock = lockManager.acquireReadLock()) { Preconditions.checkArgument(nativeDatasetHandle != 0, "Dataset is closed"); - // Flattened as [id0, rowCount0, dataFileNum0, id1, ...] to keep the JNI surface primitive - long[] flat = nativeGetFragmentStatistics(); - int count = flat.length / 3; - int[] ids = new int[count]; - long[] rowCounts = new long[count]; - int[] dataFileNums = new int[count]; - for (int i = 0; i < count; i++) { - ids[i] = (int) flat[3 * i]; - rowCounts[i] = flat[3 * i + 1]; - dataFileNums[i] = (int) flat[3 * i + 2]; - } - return new FragmentStatistics(ids, rowCounts, dataFileNums); - } - } - - private native long[] nativeGetFragmentStatistics(); - - /** - * Get aggregate statistics for all fragments in this dataset version. - * - *

The aggregation runs in native code and returns a fixed-size result, avoiding per-fragment - * Java objects and arrays. - * - *

Every fragment must contain enough metadata to determine its live row count. Some legacy - * datasets do not contain this metadata and are not supported by this method. - * - * @return aggregate fragment statistics - * @throws RuntimeException if a fragment is missing required row-count metadata - */ - public FragmentSummary getFragmentSummary() { - try (LockManager.ReadLock readLock = lockManager.acquireReadLock()) { - Preconditions.checkArgument(nativeDatasetHandle != 0, "Dataset is closed"); - return nativeGetFragmentSummary(); + return nativeGetFragmentStatistics(); } } - private native FragmentSummary nativeGetFragmentSummary(); + private native FragmentStatistics nativeGetFragmentStatistics(); /** * Gets the arrow schema of the dataset. diff --git a/java/src/main/java/org/lance/FragmentSummary.java b/java/src/main/java/org/lance/FragmentSummary.java deleted file mode 100644 index c755f2ef424..00000000000 --- a/java/src/main/java/org/lance/FragmentSummary.java +++ /dev/null @@ -1,61 +0,0 @@ -/* - * Licensed 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. - */ -package org.lance; - -/** Aggregate statistics for the fragments in a dataset version. */ -public final class FragmentSummary { - private final long fragmentCount; - private final long minRowsPerFragment; - private final long maxRowsPerFragment; - private final long minDataFilesPerFragment; - private final long maxDataFilesPerFragment; - - FragmentSummary( - long fragmentCount, - long minRowsPerFragment, - long maxRowsPerFragment, - long minDataFilesPerFragment, - long maxDataFilesPerFragment) { - this.fragmentCount = fragmentCount; - this.minRowsPerFragment = minRowsPerFragment; - this.maxRowsPerFragment = maxRowsPerFragment; - this.minDataFilesPerFragment = minDataFilesPerFragment; - this.maxDataFilesPerFragment = maxDataFilesPerFragment; - } - - /** Number of fragments. */ - public long getFragmentCount() { - return fragmentCount; - } - - /** Minimum number of live rows in a fragment, or 0 when there are no fragments. */ - public long getMinRowsPerFragment() { - return minRowsPerFragment; - } - - /** Maximum number of live rows in a fragment, or 0 when there are no fragments. */ - public long getMaxRowsPerFragment() { - return maxRowsPerFragment; - } - - /** Minimum number of data files in a fragment, or 0 when there are no fragments. */ - public long getMinDataFilesPerFragment() { - return minDataFilesPerFragment; - } - - /** Maximum number of data files in a fragment, or 0 when there are no fragments. */ - public long getMaxDataFilesPerFragment() { - return maxDataFilesPerFragment; - } -} diff --git a/java/src/test/java/org/lance/FragmentTest.java b/java/src/test/java/org/lance/FragmentTest.java index cecab4e6197..3db9cb22db7 100644 --- a/java/src/test/java/org/lance/FragmentTest.java +++ b/java/src/test/java/org/lance/FragmentTest.java @@ -442,12 +442,8 @@ void testFragmentStatistics(@TempDir Path tempDir) { assertEquals(30, Arrays.stream(stats.getRowCounts()).sum()); - FragmentSummary summary = dataset.getFragmentSummary(); - assertEquals(2, summary.getFragmentCount()); - assertEquals(9, summary.getMinRowsPerFragment()); - assertEquals(21, summary.getMaxRowsPerFragment()); - assertEquals(1, summary.getMinDataFilesPerFragment()); - assertEquals(1, summary.getMaxDataFilesPerFragment()); + dataset.delete("id < 5"); + assertArrayEquals(new long[] {16, 4}, dataset.getFragmentStatistics().getRowCounts()); } } } @@ -460,18 +456,48 @@ void testFragmentStatisticsOnEmptyDataset(@TempDir Path tempDir) { new TestUtils.SimpleTestDataset(allocator, datasetPath); try (Dataset dataset = testDataset.createEmptyDataset()) { assertEquals(0, dataset.getFragmentStatistics().size()); - FragmentSummary summary = dataset.getFragmentSummary(); - assertEquals(0, summary.getFragmentCount()); - assertEquals(0, summary.getMinRowsPerFragment()); - assertEquals(0, summary.getMaxRowsPerFragment()); - assertEquals(0, summary.getMinDataFilesPerFragment()); - assertEquals(0, summary.getMaxDataFilesPerFragment()); } } } @Test - void testFragmentSummaryRejectsUnknownRowsFromHistoricalManifest() { + void testFragmentStatisticsAcrossNativeChunks(@TempDir Path tempDir) { + String datasetPath = tempDir.resolve("fragment_statistics_chunks").toString(); + try (RootAllocator allocator = new RootAllocator(Long.MAX_VALUE)) { + TestUtils.SimpleTestDataset testDataset = + new TestUtils.SimpleTestDataset(allocator, datasetPath); + testDataset.createEmptyDataset().close(); + + FragmentMetadata template = testDataset.createNewFragment(1); + int fragmentCount = 4097; + List fragments = new ArrayList<>(fragmentCount); + for (int id = 0; id < fragmentCount; id++) { + fragments.add( + new FragmentMetadata( + id, + template.getFiles(), + template.getPhysicalRows(), + template.getDeletionFile(), + template.getRowIdMeta())); + } + + FragmentOperation.Append appendOp = new FragmentOperation.Append(fragments); + try (Dataset dataset = Dataset.commit(allocator, datasetPath, appendOp, Optional.of(1L))) { + FragmentStatistics stats = dataset.getFragmentStatistics(); + int lastIndex = fragmentCount - 1; + assertEquals(fragmentCount, stats.size()); + assertEquals(0, stats.getIds()[0]); + assertEquals(lastIndex, stats.getIds()[lastIndex]); + assertEquals(1, stats.getRowCounts()[0]); + assertEquals(1, stats.getRowCounts()[lastIndex]); + assertEquals(1, stats.getDataFileNums()[0]); + assertEquals(1, stats.getDataFileNums()[lastIndex]); + } + } + } + + @Test + void testFragmentStatisticsPreservesLegacyMissingRowCount() { String historicalPath = Path.of("..", "test_data", "v0.7.5", "with_deletions") .toAbsolutePath() @@ -479,9 +505,10 @@ void testFragmentSummaryRejectsUnknownRowsFromHistoricalManifest() { .toString(); try (RootAllocator allocator = new RootAllocator(Long.MAX_VALUE); Dataset dataset = Dataset.open(historicalPath, allocator)) { - RuntimeException error = assertThrows(RuntimeException.class, dataset::getFragmentSummary); - assertTrue(error.getMessage().contains("Fragment summary requires")); - assertTrue(error.getMessage().contains("fragment")); + FragmentStatistics stats = dataset.getFragmentStatistics(); + assertArrayEquals(new int[] {0}, stats.getIds()); + assertArrayEquals(new long[] {0}, stats.getRowCounts()); + assertArrayEquals(new int[] {1}, stats.getDataFileNums()); } } } diff --git a/rust/lance/Cargo.toml b/rust/lance/Cargo.toml index 180bb82363c..c168803e8f4 100644 --- a/rust/lance/Cargo.toml +++ b/rust/lance/Cargo.toml @@ -208,7 +208,7 @@ name = "count_pushdown" harness = false [[bench]] -name = "fragment_summary" +name = "fragment_statistics" harness = false [[bench]] diff --git a/rust/lance/benches/fragment_statistics.rs b/rust/lance/benches/fragment_statistics.rs new file mode 100644 index 00000000000..d6d148525c1 --- /dev/null +++ b/rust/lance/benches/fragment_statistics.rs @@ -0,0 +1,152 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +//! Compares the old and new per-fragment JNI payload preparation paths. +//! +//! The flattened baseline mirrors the old Java `Dataset.getFragmentStatistics()` path: it creates +//! `FileFragment` wrappers and allocates three `i64` values per fragment before JNI copies and +//! Java-side array splitting. The typed-chunk path mirrors the current implementation's native +//! preparation before it copies directly into the three final Java arrays. +//! +//! At 100,000 fragments, the old flattened path allocates a 2.4 MB native vector, followed by a +//! 2.4 MB Java `long[]` JNI copy and 1.6 MB across the three final Java primitive arrays. The +//! typed-chunk path bounds native staging memory to 64 KiB and creates only the 1.6 MB final Java +//! arrays. +//! +//! ```text +//! cargo bench -p lance --bench fragment_statistics +//! ``` + +use std::hint::black_box; + +use arrow_schema::{DataType, Field as ArrowField, Schema as ArrowSchema}; +use criterion::{Criterion, Throughput, criterion_group, criterion_main}; +use lance::Dataset; +use lance::dataset::transaction::Operation; +use lance_core::utils::tempfile::TempStrDir; +use lance_table::format::Fragment; + +const NUM_FRAGMENTS: usize = 100_000; +const STATISTICS_CHUNK_SIZE: usize = 4096; + +struct Fixture { + _data_dir: TempStrDir, + dataset: Dataset, +} + +impl Fixture { + async fn open() -> Self { + let data_dir = TempStrDir::default(); + let schema = + lance_core::datatypes::Schema::try_from(&ArrowSchema::new(vec![ArrowField::new( + "value", + DataType::Int64, + false, + )])) + .unwrap(); + let fragments = (0..NUM_FRAGMENTS) + .map(|id| { + let mut fragment = Fragment::new(id as u64); + fragment.physical_rows = Some(1_000 + id % 100); + fragment + }) + .collect(); + let operation = Operation::Overwrite { + fragments, + schema, + config_upsert_values: None, + initial_bases: None, + }; + let dataset = Dataset::commit( + data_dir.as_str(), + operation, + None, + None, + None, + Default::default(), + false, + ) + .await + .unwrap(); + + Self { + _data_dir: data_dir, + dataset, + } + } +} + +fn legacy_flattened_fragment_statistics(dataset: &Dataset) -> Vec { + let fragments = dataset.get_fragments(); + let mut statistics = Vec::with_capacity(fragments.len() * 3); + for fragment in fragments.iter() { + let metadata = fragment.metadata(); + let physical_rows = metadata.physical_rows.unwrap_or(0) as i64; + let deleted_rows = metadata + .deletion_file + .as_ref() + .and_then(|deletion_file| deletion_file.num_deleted_rows) + .unwrap_or(0) as i64; + statistics.push(metadata.id as i64); + statistics.push(physical_rows - deleted_rows); + statistics.push(metadata.files.len() as i64); + } + statistics +} + +fn prepare_typed_fragment_statistics_chunks(dataset: &Dataset) -> usize { + let fragments = dataset.fragments(); + let chunk_capacity = fragments.len().min(STATISTICS_CHUNK_SIZE); + let mut ids = Vec::with_capacity(chunk_capacity); + let mut row_counts = Vec::with_capacity(chunk_capacity); + let mut data_file_nums = Vec::with_capacity(chunk_capacity); + let mut value_count = 0; + + for fragments in fragments.chunks(STATISTICS_CHUNK_SIZE) { + ids.clear(); + row_counts.clear(); + data_file_nums.clear(); + for fragment in fragments { + let physical_rows = fragment.physical_rows.unwrap_or(0) as i64; + let deleted_rows = fragment + .deletion_file + .as_ref() + .and_then(|deletion_file| deletion_file.num_deleted_rows) + .unwrap_or(0) as i64; + ids.push(fragment.id as i32); + row_counts.push(physical_rows - deleted_rows); + data_file_nums.push(fragment.files.len() as i32); + } + value_count += ids.len() + row_counts.len() + data_file_nums.len(); + black_box((&ids, &row_counts, &data_file_nums)); + } + + value_count +} + +fn bench_fragment_statistics(c: &mut Criterion) { + let runtime = tokio::runtime::Runtime::new().unwrap(); + let fixture = runtime.block_on(Fixture::open()); + + assert_eq!( + legacy_flattened_fragment_statistics(&fixture.dataset).len(), + NUM_FRAGMENTS * 3 + ); + assert_eq!( + prepare_typed_fragment_statistics_chunks(&fixture.dataset), + NUM_FRAGMENTS * 3 + ); + + let mut group = c.benchmark_group("fragment_statistics/100k_fragments"); + group.throughput(Throughput::Elements(NUM_FRAGMENTS as u64)); + group.bench_function("legacy_materialize_flattened_statistics", |b| { + b.iter(|| black_box(legacy_flattened_fragment_statistics(&fixture.dataset))) + }); + group.bench_function("prepare_typed_statistics_chunks", |b| { + b.iter(|| black_box(prepare_typed_fragment_statistics_chunks(&fixture.dataset))) + }); + group.finish(); +} + +criterion_group!(benches, bench_fragment_statistics); +criterion_main!(benches); diff --git a/rust/lance/benches/fragment_summary.rs b/rust/lance/benches/fragment_summary.rs deleted file mode 100644 index 8fa45d99389..00000000000 --- a/rust/lance/benches/fragment_summary.rs +++ /dev/null @@ -1,117 +0,0 @@ -// SPDX-License-Identifier: Apache-2.0 -// SPDX-FileCopyrightText: Copyright The Lance Authors - -//! Compares fixed-size fragment aggregation with materializing per-fragment JNI payload data. -//! -//! The flattened baseline mirrors the native work behind Java -//! `Dataset.getFragmentStatistics()`: it allocates three `i64` values per fragment before JNI -//! copies and Java-side array splitting. `Dataset.fragment_summary()` returns five scalars -//! regardless of fragment count. -//! -//! At 100,000 fragments, the flattened path allocates a 2.4 MB native vector, followed by a -//! 2.4 MB Java `long[]` JNI copy and 1.6 MB across the three final Java primitive arrays. The -//! aggregate path creates one fixed-size Java object. -//! -//! ```text -//! cargo bench -p lance --bench fragment_summary -//! ``` - -use std::hint::black_box; - -use arrow_schema::{DataType, Field as ArrowField, Schema as ArrowSchema}; -use criterion::{Criterion, Throughput, criterion_group, criterion_main}; -use lance::Dataset; -use lance::dataset::transaction::Operation; -use lance_core::utils::tempfile::TempStrDir; -use lance_table::format::Fragment; - -const NUM_FRAGMENTS: usize = 100_000; - -struct Fixture { - _data_dir: TempStrDir, - dataset: Dataset, -} - -impl Fixture { - async fn open() -> Self { - let data_dir = TempStrDir::default(); - let schema = - lance_core::datatypes::Schema::try_from(&ArrowSchema::new(vec![ArrowField::new( - "value", - DataType::Int64, - false, - )])) - .unwrap(); - let fragments = (0..NUM_FRAGMENTS) - .map(|id| { - let mut fragment = Fragment::new(id as u64); - fragment.physical_rows = Some(1_000 + id % 100); - fragment - }) - .collect(); - let operation = Operation::Overwrite { - fragments, - schema, - config_upsert_values: None, - initial_bases: None, - }; - let dataset = Dataset::commit( - data_dir.as_str(), - operation, - None, - None, - None, - Default::default(), - false, - ) - .await - .unwrap(); - - Self { - _data_dir: data_dir, - dataset, - } - } -} - -fn flattened_fragment_statistics(dataset: &Dataset) -> Vec { - let fragments = dataset.fragments(); - let mut statistics = Vec::with_capacity(fragments.len() * 3); - for fragment in fragments.iter() { - let physical_rows = fragment.physical_rows.unwrap_or(0) as i64; - let deleted_rows = fragment - .deletion_file - .as_ref() - .and_then(|deletion_file| deletion_file.num_deleted_rows) - .unwrap_or(0) as i64; - statistics.push(fragment.id as i64); - statistics.push(physical_rows - deleted_rows); - statistics.push(fragment.files.len() as i64); - } - statistics -} - -fn bench_fragment_statistics(c: &mut Criterion) { - let runtime = tokio::runtime::Runtime::new().unwrap(); - let fixture = runtime.block_on(Fixture::open()); - - let summary = fixture.dataset.fragment_summary().unwrap(); - assert_eq!(summary.fragment_count, NUM_FRAGMENTS as u64); - assert_eq!( - flattened_fragment_statistics(&fixture.dataset).len(), - NUM_FRAGMENTS * 3 - ); - - let mut group = c.benchmark_group("fragment_statistics/100k_fragments"); - group.throughput(Throughput::Elements(NUM_FRAGMENTS as u64)); - group.bench_function("aggregate_summary", |b| { - b.iter(|| black_box(fixture.dataset.fragment_summary().unwrap())) - }); - group.bench_function("materialize_flattened_statistics", |b| { - b.iter(|| black_box(flattened_fragment_statistics(&fixture.dataset))) - }); - group.finish(); -} - -criterion_group!(benches, bench_fragment_statistics); -criterion_main!(benches); diff --git a/rust/lance/src/dataset/statistics.rs b/rust/lance/src/dataset/statistics.rs index f5ff650127a..8231a13b30f 100644 --- a/rust/lance/src/dataset/statistics.rs +++ b/rust/lance/src/dataset/statistics.rs @@ -17,141 +17,6 @@ use super::overlay::{collect_overlay_stale_frags, overlaid_fragments}; use super::{Dataset, fragment::FileFragment, versions}; use crate::index::{DatasetIndexExt, DatasetIndexInternalExt}; -/// Aggregate statistics for the fragments in a dataset version. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub struct FragmentSummary { - /// Number of fragments. - pub fragment_count: u64, - /// Minimum number of live rows in a fragment, or 0 when the dataset has no fragments. - pub min_rows_per_fragment: u64, - /// Maximum number of live rows in a fragment, or 0 when the dataset has no fragments. - pub max_rows_per_fragment: u64, - /// Minimum number of data files in a fragment, or 0 when the dataset has no fragments. - pub min_data_files_per_fragment: u64, - /// Maximum number of data files in a fragment, or 0 when the dataset has no fragments. - pub max_data_files_per_fragment: u64, -} - -impl Dataset { - /// Aggregate fragment statistics from the loaded manifest in one pass. - /// - /// Returns an error for legacy fragments that do not contain enough metadata to determine - /// their live row count. - pub fn fragment_summary(&self) -> Result { - summarize_fragments(self.fragments()) - } -} - -fn summarize_fragments(fragments: &[lance_table::format::Fragment]) -> Result { - let mut min_rows_per_fragment = u64::MAX; - let mut max_rows_per_fragment = 0; - let mut min_data_files_per_fragment = u64::MAX; - let mut max_data_files_per_fragment = 0; - - for fragment in fragments { - let live_rows = fragment.num_rows().ok_or_else(|| { - Error::internal(format!( - "Fragment summary requires physical row count and deletion count in fragment metadata, but fragment {} is missing required row-count metadata. Rewrite the dataset with a current Lance version to populate it", - fragment.id - )) - })? as u64; - min_rows_per_fragment = min_rows_per_fragment.min(live_rows); - max_rows_per_fragment = max_rows_per_fragment.max(live_rows); - let data_file_count = fragment.files.len() as u64; - min_data_files_per_fragment = min_data_files_per_fragment.min(data_file_count); - max_data_files_per_fragment = max_data_files_per_fragment.max(data_file_count); - } - - if fragments.is_empty() { - min_rows_per_fragment = 0; - min_data_files_per_fragment = 0; - } - - Ok(FragmentSummary { - fragment_count: fragments.len() as u64, - min_rows_per_fragment, - max_rows_per_fragment, - min_data_files_per_fragment, - max_data_files_per_fragment, - }) -} - -#[cfg(test)] -mod fragment_summary_tests { - use lance_core::Error; - use lance_file::version::ConcreteFileVersion; - use lance_table::format::{DeletionFile, DeletionFileType, Fragment}; - - use super::summarize_fragments; - - fn fragment(id: u64, rows: Option, file_count: usize) -> Fragment { - let mut fragment = Fragment::new(id); - fragment.physical_rows = rows; - for file_idx in 0..file_count { - fragment.add_file( - format!("{id}-{file_idx}.lance"), - vec![file_idx as i32], - vec![file_idx as i32], - ConcreteFileVersion::V2_1, - None, - ); - } - fragment - } - - #[test] - fn test_fragment_summary_rejects_unknown_row_counts() { - let known = fragment(0, Some(10), 1); - let unknown_physical_rows = fragment(1, None, 2); - let mut unknown_deletions = fragment(2, Some(30), 3); - unknown_deletions.deletion_file = Some(DeletionFile { - read_version: 1, - id: 1, - file_type: DeletionFileType::Array, - num_deleted_rows: None, - base_id: None, - }); - - let physical_rows_error = - summarize_fragments(&[known.clone(), unknown_physical_rows]).unwrap_err(); - assert!(matches!(&physical_rows_error, Error::Internal { .. })); - assert!(physical_rows_error.to_string().contains("fragment 1")); - - let deletion_count_error = summarize_fragments(&[known, unknown_deletions]).unwrap_err(); - assert!(matches!(&deletion_count_error, Error::Internal { .. })); - assert!(deletion_count_error.to_string().contains("fragment 2")); - } - - #[test] - fn test_fragment_summary_uses_live_rows() { - let mut deleted = fragment(0, Some(10), 1); - deleted.deletion_file = Some(DeletionFile { - read_version: 1, - id: 1, - file_type: DeletionFileType::Array, - num_deleted_rows: Some(4), - base_id: None, - }); - let other = fragment(1, Some(20), 2); - - let summary = summarize_fragments(&[deleted, other]).unwrap(); - assert_eq!(summary.min_rows_per_fragment, 6); - assert_eq!(summary.max_rows_per_fragment, 20); - assert_eq!(summary.min_data_files_per_fragment, 1); - assert_eq!(summary.max_data_files_per_fragment, 2); - } - - #[test] - fn test_fragment_summary_empty() { - let summary = summarize_fragments(&[]).unwrap(); - assert_eq!(summary.fragment_count, 0); - assert_eq!(summary.min_rows_per_fragment, 0); - assert_eq!(summary.max_rows_per_fragment, 0); - assert_eq!(summary.min_data_files_per_fragment, 0); - assert_eq!(summary.max_data_files_per_fragment, 0); - } -} - /// Statistics about a single field in the dataset pub struct FieldStatistics { /// Id of the field From 7995e57a7e592a0ed1709deb9ffb9fb55aec128a Mon Sep 17 00:00:00 2001 From: wangzheyan Date: Fri, 14 Aug 2026 00:55:41 +0800 Subject: [PATCH 6/6] ci: rerun checks