diff --git a/crates/paimon/src/spec/manifest_sidecar.rs b/crates/paimon/src/spec/manifest_sidecar.rs index c11a4a41e..ecdeafe6e 100644 --- a/crates/paimon/src/spec/manifest_sidecar.rs +++ b/crates/paimon/src/spec/manifest_sidecar.rs @@ -44,6 +44,8 @@ const MIN_BLOCK_BYTES: usize = 6; const MAX_INT: u64 = i32::MAX as u64; const MAX_LONG: u64 = i64::MAX as u64; const BLOCK_READ_BUFFER_BYTES: u64 = 4 * 1024 * 1024; +// Cap the extra bytes spent to replace several range requests with one full read. +const MAX_FULL_READ_AMPLIFICATION: usize = 4; type RowBlockFilter<'a> = dyn Fn(i64, i64) -> bool + Sync + 'a; type PartitionBlockFilter<'a> = dyn FnMut(&[u8]) -> bool + Send + 'a; @@ -81,6 +83,7 @@ impl ManifestSidecarBlock { pub struct ManifestSidecarSelection { header: Bytes, blocks: Vec, + manifest_size: u64, } impl ManifestSidecarSelection { @@ -91,6 +94,64 @@ impl ManifestSidecarSelection { pub fn blocks(&self) -> &[ManifestSidecarBlock] { &self.blocks } + + /// Use a full read when it saves requests without multiplying transferred + /// bytes by more than the number of range reads (capped at four). + fn should_read_full(&self) -> Result { + let range_reads = self.read_ranges()?.len(); + if range_reads <= 1 { + return Ok(false); + } + let selected_bytes = self.blocks.iter().try_fold(0u64, |total, block| { + total.checked_add(block.length).ok_or_else(invalid_sidecar) + })?; + Ok( + u128::from(selected_bytes) * range_reads.min(MAX_FULL_READ_AMPLIFICATION) as u128 + >= u128::from(self.manifest_size), + ) + } + + pub(crate) async fn read_bytes(&self, file_io: &FileIO, manifest_path: &str) -> Result { + if self.should_read_full()? { + file_io.new_input(manifest_path)?.read().await + } else { + ManifestSidecar::read_selected_bytes(file_io, manifest_path, self).await + } + } + + fn read_ranges(&self) -> Result>> { + let mut ranges = Vec::new(); + let mut block_index = 0usize; + while block_index < self.blocks.len() { + let first = self.blocks[block_index]; + let mut span_end = first + .offset + .checked_add(first.length) + .ok_or_else(invalid_sidecar)?; + block_index += 1; + while block_index < self.blocks.len() { + let next = self.blocks[block_index]; + let Some(next_end) = next.offset.checked_add(next.length) else { + return Err(invalid_sidecar()); + }; + if next.offset != span_end + || next_end.saturating_sub(first.offset) > BLOCK_READ_BUFFER_BYTES + { + break; + } + span_end = next_end; + block_index += 1; + } + + let mut position = first.offset; + while position < span_end { + let end = span_end.min(position + BLOCK_READ_BUFFER_BYTES); + ranges.push(position..end); + position = end; + } + } + Ok(ranges) + } } #[derive(Debug)] @@ -566,41 +627,10 @@ impl ManifestSidecar { ); out.extend_from_slice(&selection.header); - let mut block_index = 0usize; - while block_index < selection.blocks.len() { - let first = selection.blocks[block_index]; - let mut span_end = first - .offset - .checked_add(first.length) - .ok_or_else(invalid_sidecar)?; - block_index += 1; - while block_index < selection.blocks.len() { - let next = selection.blocks[block_index]; - let Some(next_end) = next.offset.checked_add(next.length) else { - return Err(invalid_sidecar()); - }; - if next.offset != span_end - || next_end.saturating_sub(first.offset) > BLOCK_READ_BUFFER_BYTES - { - break; - } - span_end = next_end; - block_index += 1; - } - - let mut position = first.offset; - while position < span_end { - let end = span_end.min(position + BLOCK_READ_BUFFER_BYTES); - let bytes = reader - .read(Range { - start: position, - end, - }) - .await?; - require(bytes.len() as u64 == end - position)?; - out.extend_from_slice(&bytes); - position = end; - } + for range in selection.read_ranges()? { + let bytes = reader.read(range.clone()).await?; + require(bytes.len() as u64 == range.end - range.start)?; + out.extend_from_slice(&bytes); } Ok(Bytes::from(out)) } @@ -846,6 +876,7 @@ fn select_with_filters( Ok(ManifestSidecarSelection { header: Bytes::copy_from_slice(header), blocks: selected, + manifest_size, }) } @@ -1079,6 +1110,54 @@ mod tests { (builder.serialize(size, 7).unwrap(), meta(size, 7)) } + #[test] + fn full_read_policy_accounts_for_bytes_and_range_requests() { + const MIB: u64 = 1024 * 1024; + let selection = |blocks: &[(u64, u64)]| ManifestSidecarSelection { + header: Bytes::new(), + blocks: blocks + .iter() + .map(|&(offset, length)| ManifestSidecarBlock { + offset, + length, + first_record: 0, + record_count: 1, + }) + .collect(), + manifest_size: 8 * MIB, + }; + + assert!(!selection(&[]).should_read_full().unwrap()); + // A contiguous selection larger than the 4 MiB read buffer still needs + // two GETs, so a full read saves a request. + let large = selection(&[(0, 6 * MIB)]); + assert_eq!(large.read_ranges().unwrap().len(), 2); + assert!(large.should_read_full().unwrap()); + + assert!(selection(&[(0, 2 * MIB), (4 * MIB, 2 * MIB)]) + .should_read_full() + .unwrap()); + assert!(!selection(&[(0, MIB), (4 * MIB, MIB)]) + .should_read_full() + .unwrap()); + assert!(selection(&[ + (0, MIB / 2), + (2 * MIB, MIB / 2), + (4 * MIB, MIB / 2), + (6 * MIB, MIB / 2), + ]) + .should_read_full() + .unwrap()); + assert!(!selection(&[ + (0, MIB / 8), + (2 * MIB, MIB / 8), + (4 * MIB, MIB / 8), + (6 * MIB, MIB / 8), + ]) + .should_read_full() + .unwrap()); + } + fn partition(p: i32, q: Option<&str>) -> Vec { let mut builder = BinaryRowBuilder::new(2); builder.write_int(0, p); diff --git a/crates/paimon/src/table/table_scan.rs b/crates/paimon/src/table/table_scan.rs index f4edba05b..e2a260328 100644 --- a/crates/paimon/src/table/table_scan.rs +++ b/crates/paimon/src/table/table_scan.rs @@ -189,7 +189,7 @@ async fn read_manifest_bytes_with_sidecar( .await; match selection.as_ref() { - Some(selection) => ManifestSidecar::read_selected_bytes(file_io, path, selection).await, + Some(selection) => selection.read_bytes(file_io, path).await, None => file_io.new_input(path)?.read().await, } } @@ -2746,11 +2746,11 @@ mod tests { use super::{ data_evolution_row_range_groups, data_file_overlaps_row_range_index, group_data_files_by_partition_bucket, manifest_file_overlaps_row_range_index, - prune_data_evolution_group_by_read_fields, retain_index_manifest_entry, - retain_index_manifest_entry_for_scan, retain_manifest_buckets, + prune_data_evolution_group_by_read_fields, read_manifest_bytes_with_sidecar, + retain_index_manifest_entry, retain_index_manifest_entry_for_scan, retain_manifest_buckets, retain_manifest_entry_row_ranges, retain_manifest_row_ranges, scan_predicate_field_ids, should_skip_level_zero_for_scan, split_row_ranges_for_files, LimitPushdownAccumulator, - PaimonTableScan, RowRangeIndex, TableScan, + ManifestSidecarPruning, PaimonTableScan, RowRangeIndex, TableScan, }; use crate::catalog::Identifier; use crate::io::FileIOBuilder; @@ -2758,8 +2758,9 @@ mod tests { stats::BinaryTableStats, ArrayType, BinaryRow, BinaryRowBuilder, BucketFunctionType, ColumnMove, CommitKind, DataField, DataFileMeta, DataType, Datum, DeletionVectorMeta, FileKind, GlobalIndexMeta, IndexFileMeta, IndexManifestEntry, IntType, ManifestEntry, - ManifestFileMeta, Predicate, PredicateBuilder, PredicateOperator, Schema as PaimonSchema, - SchemaChange, Snapshot, TableSchema, VarCharType, + ManifestFileMeta, ManifestSidecar, ManifestSidecarBuilder, Predicate, PredicateBuilder, + PredicateOperator, Schema as PaimonSchema, SchemaChange, Snapshot, TableSchema, + VarCharType, }; use crate::table::bucket_filter::{compute_target_buckets, extract_predicate_for_keys}; use crate::table::partition_filter::PartitionFilter; @@ -3946,6 +3947,90 @@ mod tests { assert_eq!(delta_files, vec!["a-new", "a-old"]); } + #[tokio::test] + async fn test_manifest_sidecar_read_uses_full_file_only_for_dense_selection() { + let file_io = FileIOBuilder::new("memory").build().unwrap(); + let manifest_path = "memory:/adaptive_manifest_read/manifest-test"; + let sidecar_path = ManifestSidecar::path(manifest_path); + let mut header = b"Obj\x01".to_vec(); + header.resize(32, 0); + let mut builder = ManifestSidecarBuilder::new(header.clone(), true, false); + let mut offset = header.len() as u64; + for (length, first_row_id) in [(200, 0), (50, 100), (200, 200)] { + builder.begin_block(offset, length, 1).unwrap(); + builder + .add(Some(first_row_id), 1, None, None, None) + .unwrap(); + builder.end_block().unwrap(); + offset += length; + } + let sidecar = builder.serialize(offset, 3).unwrap(); + let manifest = ManifestFileMeta::new( + "manifest-test".to_string(), + offset as i64, + 3, + 0, + BinaryTableStats::empty(), + 0, + ) + .with_extra_files(Some(vec!["manifest-test.avro.sidecar".to_string()])); + let mut avro = header.clone(); + avro.extend((0..450).map(|value| value as u8)); + file_io + .new_output(manifest_path) + .unwrap() + .write(Bytes::from(avro.clone())) + .await + .unwrap(); + file_io + .new_output(&sidecar_path) + .unwrap() + .write(Bytes::from(sidecar)) + .await + .unwrap(); + + let dense = RowRangeIndex::create(vec![RowRange::new(0, 0), RowRange::new(200, 200)]); + let dense_bytes = read_manifest_bytes_with_sidecar( + &file_io, + manifest_path, + &manifest, + true, + ManifestSidecarPruning { + row_range_index: Some(&dense), + partition_filter: None, + partition_arity: 0, + bucket_predicate: None, + bucket_key_fields: &[], + bucket_function_type: BucketFunctionType::Default, + }, + ) + .await + .unwrap(); + assert_eq!(dense_bytes.as_ref(), avro); + + let sparse = RowRangeIndex::create(vec![RowRange::new(0, 0)]); + let sparse_bytes = read_manifest_bytes_with_sidecar( + &file_io, + manifest_path, + &manifest, + true, + ManifestSidecarPruning { + row_range_index: Some(&sparse), + partition_filter: None, + partition_arity: 0, + bucket_predicate: None, + bucket_key_fields: &[], + bucket_function_type: BucketFunctionType::Default, + }, + ) + .await + .unwrap(); + let mut expected = header; + let start = expected.len(); + expected.extend_from_slice(&avro[start..start + 200]); + assert_eq!(sparse_bytes.as_ref(), expected); + } + #[tokio::test] async fn test_manifest_sidecar_prunes_avro_blocks_before_entry_decode() { const FILE_COUNT: usize = 1_000;