Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
149 changes: 114 additions & 35 deletions crates/paimon/src/spec/manifest_sidecar.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -81,6 +83,7 @@ impl ManifestSidecarBlock {
pub struct ManifestSidecarSelection {
header: Bytes,
blocks: Vec<ManifestSidecarBlock>,
manifest_size: u64,
}

impl ManifestSidecarSelection {
Expand All @@ -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<bool> {
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<Bytes> {
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<Vec<Range<u64>>> {
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)]
Expand Down Expand Up @@ -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))
}
Expand Down Expand Up @@ -846,6 +876,7 @@ fn select_with_filters(
Ok(ManifestSidecarSelection {
header: Bytes::copy_from_slice(header),
blocks: selected,
manifest_size,
})
}

Expand Down Expand Up @@ -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<u8> {
let mut builder = BinaryRowBuilder::new(2);
builder.write_int(0, p);
Expand Down
97 changes: 91 additions & 6 deletions crates/paimon/src/table/table_scan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}
}
Expand Down Expand Up @@ -2746,20 +2746,21 @@ 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;
use crate::spec::{
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;
Expand Down Expand Up @@ -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;
Expand Down
Loading