diff --git a/crates/paimon/src/btree/block.rs b/crates/paimon/src/btree/block.rs index 6f5c513b0..6bb8ec27d 100644 --- a/crates/paimon/src/btree/block.rs +++ b/crates/paimon/src/btree/block.rs @@ -588,6 +588,16 @@ impl BlockReader { } } + pub(crate) fn retained_bytes(&self) -> usize { + let seek_bytes = match &self.seek_info { + SeekInfo::Aligned { .. } => 0, + SeekInfo::Unaligned { offsets } => offsets + .capacity() + .saturating_mul(std::mem::size_of::()), + }; + self.data.capacity().saturating_add(seek_bytes) + } + #[allow(dead_code)] pub fn record_count(&self) -> usize { self.record_count diff --git a/crates/paimon/src/btree/data_block_cache.rs b/crates/paimon/src/btree/data_block_cache.rs new file mode 100644 index 000000000..530dc0019 --- /dev/null +++ b/crates/paimon/src/btree/data_block_cache.rs @@ -0,0 +1,104 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you 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. + +use crate::btree::block::BlockReader; +use lru::LruCache; +use std::sync::{Arc, Mutex}; + +#[derive(Clone, PartialEq, Eq, Hash)] +pub(super) struct DataBlockCacheKey { + pub(super) file: Arc, + pub(super) offset: u64, + pub(super) size: u32, +} + +struct Entry { + block: Arc, + retained_bytes: usize, +} + +struct State { + entries: LruCache, + retained_bytes: usize, +} + +pub(crate) struct BTreeDataBlockCache { + max_bytes: usize, + state: Mutex, +} + +impl BTreeDataBlockCache { + pub(crate) fn new(max_bytes: usize) -> Self { + Self { + max_bytes, + state: Mutex::new(State { + entries: LruCache::unbounded(), + retained_bytes: 0, + }), + } + } + + pub(super) fn get(&self, key: &DataBlockCacheKey) -> Option> { + if self.max_bytes == 0 { + return None; + } + self.state + .lock() + .unwrap_or_else(|error| error.into_inner()) + .entries + .get(key) + .map(|entry| Arc::clone(&entry.block)) + } + + pub(super) fn put(&self, key: DataBlockCacheKey, block: Arc) -> Arc { + if self.max_bytes == 0 { + return block; + } + let retained_bytes = block.retained_bytes(); + if retained_bytes > self.max_bytes { + return block; + } + + let mut state = self.state.lock().unwrap_or_else(|error| error.into_inner()); + if let Some(entry) = state.entries.get(&key) { + return Arc::clone(&entry.block); + } + state.retained_bytes = state.retained_bytes.saturating_add(retained_bytes); + state.entries.put( + key, + Entry { + block: Arc::clone(&block), + retained_bytes, + }, + ); + while state.retained_bytes > self.max_bytes { + let Some((_, entry)) = state.entries.pop_lru() else { + break; + }; + state.retained_bytes = state.retained_bytes.saturating_sub(entry.retained_bytes); + } + block + } + + #[cfg(test)] + pub(crate) fn retained_bytes(&self) -> usize { + self.state + .lock() + .unwrap_or_else(|error| error.into_inner()) + .retained_bytes + } +} diff --git a/crates/paimon/src/btree/mod.rs b/crates/paimon/src/btree/mod.rs index 5b7fc6fab..ddee2d622 100644 --- a/crates/paimon/src/btree/mod.rs +++ b/crates/paimon/src/btree/mod.rs @@ -38,6 +38,7 @@ mod block; mod bloom_filter; +mod data_block_cache; mod footer; pub(crate) mod key_serde; mod meta; @@ -52,6 +53,7 @@ pub use block::BlockCompressionType; pub(crate) use block::{ compress_block, compress_codec_block, compute_crc32, decompress_block, decompress_codec_block, }; +pub(crate) use data_block_cache::BTreeDataBlockCache; pub use footer::BTreeFileFooter; pub use key_serde::{make_key_comparator, serialize_datum}; pub use meta::BTreeIndexMeta; diff --git a/crates/paimon/src/btree/reader.rs b/crates/paimon/src/btree/reader.rs index 4e2f50e44..2e8f30e93 100644 --- a/crates/paimon/src/btree/reader.rs +++ b/crates/paimon/src/btree/reader.rs @@ -25,6 +25,7 @@ use crate::btree::block::{BlockHandle, BlockReader}; use crate::btree::bloom_filter::BloomFilter; +use crate::btree::data_block_cache::{BTreeDataBlockCache, DataBlockCacheKey}; use crate::btree::footer::{BTreeFileFooter, BloomFilterHandle, BTREE_FOOTER_ENCODED_LENGTH}; use crate::btree::key_serde::key_comparison_io_error; use crate::btree::meta::BTreeIndexMeta; @@ -35,6 +36,7 @@ use crate::spec::murmur_hash::hash_bytes; use roaring::RoaringTreemap; use std::cmp::Ordering; use std::io; +use std::sync::Arc; use tokio::sync::OnceCell; struct LazyBloomFilter { @@ -72,6 +74,8 @@ pub struct BTreeIndexReader crate::Result> { key_comparator: F, bloom_filter: Option, file_version: u32, + data_block_cache: Arc, + data_block_cache_file: Arc, } impl crate::Result> BTreeIndexReader { @@ -83,6 +87,25 @@ impl crate::Result> BTreeIndexReader { file_size: u64, meta: &BTreeIndexMeta, key_comparator: F, + ) -> io::Result { + Self::open_with_data_block_cache( + reader, + file_size, + meta, + key_comparator, + Arc::new(BTreeDataBlockCache::new(0)), + Arc::from(""), + ) + .await + } + + pub(crate) async fn open_with_data_block_cache( + reader: Box, + file_size: u64, + meta: &BTreeIndexMeta, + key_comparator: F, + data_block_cache: Arc, + data_block_cache_file: Arc, ) -> io::Result { if file_size < BTREE_FOOTER_ENCODED_LENGTH as u64 { return Err(io::Error::new( @@ -128,6 +151,8 @@ impl crate::Result> BTreeIndexReader { key_comparator, bloom_filter, file_version: footer.version, + data_block_cache, + data_block_cache_file, }) } @@ -287,14 +312,23 @@ impl crate::Result> BTreeIndexReader { } /// Read a data block from the file on demand. - async fn read_data_block(&self, handle: &BlockHandle) -> io::Result { + async fn read_data_block(&self, handle: &BlockHandle) -> io::Result> { + let key = DataBlockCacheKey { + file: Arc::clone(&self.data_block_cache_file), + offset: handle.offset, + size: handle.size, + }; + if let Some(block) = self.data_block_cache.get(&key) { + return Ok(block); + } let end = handle.offset + handle.full_block_size() as u64; let bytes = self .reader .read(handle.offset..end) .await .map_err(|e| io::Error::other(e.to_string()))?; - read_block_from_bytes(&bytes, handle.size) + let block = Arc::new(read_block_from_bytes(&bytes, handle.size)?); + Ok(self.data_block_cache.put(key, block)) } async fn bloom_might_contain(&self, key: &[u8]) -> io::Result { diff --git a/crates/paimon/src/btree/tests.rs b/crates/paimon/src/btree/tests.rs index 105e9b71f..c9a0f20c8 100644 --- a/crates/paimon/src/btree/tests.rs +++ b/crates/paimon/src/btree/tests.rs @@ -23,6 +23,7 @@ use crate::btree::query::IndexQuery; use crate::btree::reader::BTreeIndexReader; use crate::btree::test_util::{BytesFileRead, VecFileWrite}; use crate::btree::writer::BTreeIndexWriter; +use crate::btree::BTreeDataBlockCache; use crate::io::FileRead; use crate::spec::murmur_hash::hash_bytes; use crate::spec::{DataType, Datum, PredicateOperator, VarCharType}; @@ -823,6 +824,239 @@ async fn test_sparse_in_query_reads_only_target_blocks_once() { assert_eq!(ranges.lock().unwrap().len(), 3); } +#[tokio::test] +async fn test_repeated_equal_queries_reuse_data_block() { + let buf = VecFileWrite::new(); + let mut writer = + BTreeIndexWriter::new(Box::new(buf.clone()), 1_000_000, BlockCompressionType::None); + for i in 0..100 { + writer.write(Some(&int_key(i)), i as i64).await.unwrap(); + } + + let write_result = writer.finish().await.unwrap(); + let data = Bytes::from(buf.to_vec()); + let file_size = data.len() as u64; + + let uncached_ranges = Arc::new(Mutex::new(Vec::new())); + let uncached = BTreeIndexReader::open_with_data_block_cache( + Box::new(RecordingFileRead { + data: data.clone(), + ranges: Arc::clone(&uncached_ranges), + }), + file_size, + &write_result.meta, + int_cmp, + Arc::new(BTreeDataBlockCache::new(0)), + Arc::from("memory:/uncached.btree"), + ) + .await + .unwrap(); + uncached_ranges.lock().unwrap().clear(); + for key in [1, 2, 3] { + uncached.query_equal(&int_key(key)).await.unwrap(); + } + assert_eq!(uncached_ranges.lock().unwrap().len(), 3); + + let ranges = Arc::new(Mutex::new(Vec::new())); + let reader = BTreeIndexReader::open_with_data_block_cache( + Box::new(RecordingFileRead { + data, + ranges: Arc::clone(&ranges), + }), + file_size, + &write_result.meta, + int_cmp, + Arc::new(BTreeDataBlockCache::new(1024 * 1024)), + Arc::from("memory:/first.btree"), + ) + .await + .unwrap(); + ranges.lock().unwrap().clear(); + + for key in [1, 2, 3] { + let rows = reader.query_equal(&int_key(key)).await.unwrap(); + assert_eq!(rows.iter().collect::>(), vec![key as u64]); + } + assert_eq!(ranges.lock().unwrap().len(), 1); +} + +#[tokio::test] +async fn test_data_block_cache_is_bounded_and_file_scoped() { + async fn write_file(row_id: i64) -> (Bytes, crate::btree::writer::BTreeWriteResult) { + let buf = VecFileWrite::new(); + let mut writer = + BTreeIndexWriter::new(Box::new(buf.clone()), 1_000_000, BlockCompressionType::None); + writer.write(Some(&int_key(1)), row_id).await.unwrap(); + let result = writer.finish().await.unwrap(); + (Bytes::from(buf.to_vec()), result) + } + + let cache = Arc::new(BTreeDataBlockCache::new(1024 * 1024)); + let (first_data, first_result) = write_file(1).await; + let first_ranges = Arc::new(Mutex::new(Vec::new())); + let first = BTreeIndexReader::open_with_data_block_cache( + Box::new(RecordingFileRead { + data: first_data.clone(), + ranges: Arc::clone(&first_ranges), + }), + first_data.len() as u64, + &first_result.meta, + int_cmp, + Arc::clone(&cache), + Arc::from("memory:/first.btree"), + ) + .await + .unwrap(); + + let (second_data, second_result) = write_file(2).await; + let second_ranges = Arc::new(Mutex::new(Vec::new())); + let second = BTreeIndexReader::open_with_data_block_cache( + Box::new(RecordingFileRead { + data: second_data.clone(), + ranges: Arc::clone(&second_ranges), + }), + second_data.len() as u64, + &second_result.meta, + int_cmp, + Arc::clone(&cache), + Arc::from("memory:/second.btree"), + ) + .await + .unwrap(); + first_ranges.lock().unwrap().clear(); + second_ranges.lock().unwrap().clear(); + + assert_eq!( + first + .query_equal(&int_key(1)) + .await + .unwrap() + .iter() + .collect::>(), + vec![1] + ); + assert_eq!( + second + .query_equal(&int_key(1)) + .await + .unwrap() + .iter() + .collect::>(), + vec![2] + ); + assert_eq!(first_ranges.lock().unwrap().len(), 1); + assert_eq!(second_ranges.lock().unwrap().len(), 1); + + let two_file_bytes = cache.retained_bytes(); + assert!(two_file_bytes > 1); + let bounded_cache = Arc::new(BTreeDataBlockCache::new(two_file_bytes - 1)); + let bounded_first_ranges = Arc::new(Mutex::new(Vec::new())); + let bounded_first = BTreeIndexReader::open_with_data_block_cache( + Box::new(RecordingFileRead { + data: first_data.clone(), + ranges: Arc::clone(&bounded_first_ranges), + }), + first_data.len() as u64, + &first_result.meta, + int_cmp, + Arc::clone(&bounded_cache), + Arc::from("memory:/bounded-first.btree"), + ) + .await + .unwrap(); + let bounded_second_ranges = Arc::new(Mutex::new(Vec::new())); + let bounded_second = BTreeIndexReader::open_with_data_block_cache( + Box::new(RecordingFileRead { + data: second_data.clone(), + ranges: Arc::clone(&bounded_second_ranges), + }), + second_data.len() as u64, + &second_result.meta, + int_cmp, + Arc::clone(&bounded_cache), + Arc::from("memory:/bounded-second.btree"), + ) + .await + .unwrap(); + bounded_first_ranges.lock().unwrap().clear(); + bounded_second_ranges.lock().unwrap().clear(); + + bounded_first.query_equal(&int_key(1)).await.unwrap(); + bounded_second.query_equal(&int_key(1)).await.unwrap(); + bounded_first.query_equal(&int_key(1)).await.unwrap(); + assert_eq!(bounded_first_ranges.lock().unwrap().len(), 2); + assert_eq!(bounded_second_ranges.lock().unwrap().len(), 1); + assert!(bounded_cache.retained_bytes() < two_file_bytes); + + let ranges = Arc::new(Mutex::new(Vec::new())); + let uncached = BTreeIndexReader::open_with_data_block_cache( + Box::new(RecordingFileRead { + data: first_data.clone(), + ranges: Arc::clone(&ranges), + }), + first_data.len() as u64, + &first_result.meta, + int_cmp, + Arc::new(BTreeDataBlockCache::new(0)), + Arc::from("memory:/uncached.btree"), + ) + .await + .unwrap(); + ranges.lock().unwrap().clear(); + uncached.query_equal(&int_key(1)).await.unwrap(); + uncached.query_equal(&int_key(1)).await.unwrap(); + assert_eq!(ranges.lock().unwrap().len(), 2); +} + +#[tokio::test] +async fn test_data_block_cache_evicts_by_decoded_bytes() { + let buf = VecFileWrite::new(); + let mut writer = BTreeIndexWriter::new(Box::new(buf.clone()), 64, BlockCompressionType::None); + for i in 0..100 { + writer.write(Some(&int_key(i)), i as i64).await.unwrap(); + } + let write_result = writer.finish().await.unwrap(); + let data = Bytes::from(buf.to_vec()); + let file_size = data.len() as u64; + + let measuring_cache = Arc::new(BTreeDataBlockCache::new(usize::MAX)); + let measuring_reader = BTreeIndexReader::open_with_data_block_cache( + Box::new(BytesFileRead(data.clone())), + file_size, + &write_result.meta, + int_cmp, + Arc::clone(&measuring_cache), + Arc::from("memory:/measuring.btree"), + ) + .await + .unwrap(); + measuring_reader.query_equal(&int_key(0)).await.unwrap(); + measuring_reader.query_equal(&int_key(99)).await.unwrap(); + let two_block_bytes = measuring_cache.retained_bytes(); + assert!(two_block_bytes > 2); + + let ranges = Arc::new(Mutex::new(Vec::new())); + let reader = BTreeIndexReader::open_with_data_block_cache( + Box::new(RecordingFileRead { + data, + ranges: Arc::clone(&ranges), + }), + file_size, + &write_result.meta, + int_cmp, + Arc::new(BTreeDataBlockCache::new(two_block_bytes - 1)), + Arc::from("memory:/bounded.btree"), + ) + .await + .unwrap(); + ranges.lock().unwrap().clear(); + + reader.query_equal(&int_key(0)).await.unwrap(); + reader.query_equal(&int_key(99)).await.unwrap(); + reader.query_equal(&int_key(0)).await.unwrap(); + assert_eq!(ranges.lock().unwrap().len(), 3); +} + #[tokio::test] async fn test_sparse_in_query_seeks_each_key_within_target_block() { let buf = VecFileWrite::new(); diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index 9168d4414..b159c389e 100644 --- a/crates/paimon/src/spec/core_options.rs +++ b/crates/paimon/src/spec/core_options.rs @@ -37,6 +37,8 @@ const GLOBAL_INDEX_COLUMN_UPDATE_ACTION_OPTION: &str = "global-index.column-upda pub(crate) const INDEX_FILE_IN_DATA_FILE_DIR_OPTION: &str = "index-file-in-data-file-dir"; const SORTED_INDEX_RECORDS_PER_RANGE_OPTION: &str = "sorted-index.records-per-range"; const BTREE_INDEX_RECORDS_PER_RANGE_OPTION: &str = "btree-index.records-per-range"; +const BTREE_INDEX_CACHE_SIZE_OPTION: &str = "btree-index.cache-size"; +const BTREE_INDEX_HIGH_PRIORITY_POOL_RATIO_OPTION: &str = "btree-index.high-priority-pool-ratio"; const BTREE_INDEX_FALLBACK_SCAN_MAX_SIZE_OPTION: &str = "btree-index.fallback-scan-max-size"; const BITMAP_INDEX_FALLBACK_SCAN_MAX_SIZE_OPTION: &str = "bitmap-index.fallback-scan-max-size"; const SOURCE_SPLIT_TARGET_SIZE_OPTION: &str = "source.split.target-size"; @@ -169,6 +171,8 @@ const MAX_GLOBAL_INDEX_THREAD_NUM: i64 = { } }; const DEFAULT_GLOBAL_INDEX_FALLBACK_SCAN_MAX_SIZE: i64 = 256 * 1024 * 1024; +const DEFAULT_BTREE_INDEX_CACHE_SIZE: i64 = 128 * 1024 * 1024; +const DEFAULT_BTREE_INDEX_HIGH_PRIORITY_POOL_RATIO: f64 = 0.1; const BLOB_AS_DESCRIPTOR_OPTION: &str = "blob-as-descriptor"; pub(crate) const BLOB_FIELD_OPTION: &str = "blob-field"; pub(crate) const BLOB_DESCRIPTOR_FIELD_OPTION: &str = "blob-descriptor-field"; @@ -1014,6 +1018,70 @@ impl<'a> CoreOptions<'a> { self.fallback_scan_max_size(BTREE_INDEX_FALLBACK_SCAN_MAX_SIZE_OPTION) } + pub fn btree_index_cache_size(&self) -> crate::Result { + let value = match self.options.get(BTREE_INDEX_CACHE_SIZE_OPTION) { + Some(raw) => parse_memory_size(raw).ok_or_else(|| crate::Error::DataInvalid { + message: format!( + "Option '{}' must be a valid memory size, got: {}", + BTREE_INDEX_CACHE_SIZE_OPTION, raw + ), + source: None, + })?, + None => DEFAULT_BTREE_INDEX_CACHE_SIZE, + }; + if value < 0 { + return Err(crate::Error::DataInvalid { + message: format!( + "Option '{}' must be greater than or equal to 0, got: {}", + BTREE_INDEX_CACHE_SIZE_OPTION, value + ), + source: None, + }); + } + usize::try_from(value).map_err(|_| crate::Error::DataInvalid { + message: format!( + "Option '{}' is too large: {}", + BTREE_INDEX_CACHE_SIZE_OPTION, value + ), + source: None, + }) + } + + pub fn btree_index_high_priority_pool_ratio(&self) -> crate::Result { + let value = match self + .options + .get(BTREE_INDEX_HIGH_PRIORITY_POOL_RATIO_OPTION) + { + Some(raw) => raw + .trim() + .parse::() + .map_err(|_| crate::Error::DataInvalid { + message: format!( + "Option '{}' must be a valid number, got: {}", + BTREE_INDEX_HIGH_PRIORITY_POOL_RATIO_OPTION, raw + ), + source: None, + })?, + None => DEFAULT_BTREE_INDEX_HIGH_PRIORITY_POOL_RATIO, + }; + if !(0.0..1.0).contains(&value) { + return Err(crate::Error::DataInvalid { + message: format!( + "Option '{}' must be in the range [0, 1), got: {}", + BTREE_INDEX_HIGH_PRIORITY_POOL_RATIO_OPTION, value + ), + source: None, + }); + } + Ok(value) + } + + pub fn btree_index_data_block_cache_size(&self) -> crate::Result { + let total = self.btree_index_cache_size()?; + let high_priority_ratio = self.btree_index_high_priority_pool_ratio()?; + Ok((total as f64 * (1.0 - high_priority_ratio)) as usize) + } + pub fn bitmap_index_fallback_scan_max_size(&self) -> crate::Result { self.fallback_scan_max_size(BITMAP_INDEX_FALLBACK_SCAN_MAX_SIZE_OPTION) } @@ -2456,6 +2524,74 @@ mod tests { } } + #[test] + fn test_btree_index_cache_size_matches_java_default_and_validates_values() { + assert_eq!( + CoreOptions::new(&HashMap::new()) + .btree_index_cache_size() + .unwrap(), + 128 * 1024 * 1024 + ); + + let disabled = + HashMap::from([(BTREE_INDEX_CACHE_SIZE_OPTION.to_string(), "0".to_string())]); + assert_eq!( + CoreOptions::new(&disabled) + .btree_index_cache_size() + .unwrap(), + 0 + ); + + for value in ["-1", "invalid"] { + let options = + HashMap::from([(BTREE_INDEX_CACHE_SIZE_OPTION.to_string(), value.to_string())]); + let error = CoreOptions::new(&options) + .btree_index_cache_size() + .expect_err("invalid BTree cache size should fail"); + assert!(error.to_string().contains(BTREE_INDEX_CACHE_SIZE_OPTION)); + } + } + + #[test] + fn test_btree_index_data_block_cache_size_matches_java_pool_split() { + let empty = HashMap::new(); + let defaults = CoreOptions::new(&empty); + assert_eq!( + defaults.btree_index_data_block_cache_size().unwrap(), + (128.0 * 1024.0 * 1024.0 * 0.9) as usize + ); + + let options = HashMap::from([ + ( + BTREE_INDEX_CACHE_SIZE_OPTION.to_string(), + (128 * 1024 * 1024).to_string(), + ), + ( + BTREE_INDEX_HIGH_PRIORITY_POOL_RATIO_OPTION.to_string(), + "0.9".to_string(), + ), + ]); + assert_eq!( + CoreOptions::new(&options) + .btree_index_data_block_cache_size() + .unwrap(), + (128.0 * 1024.0 * 1024.0 * 0.1) as usize + ); + + for value in ["-0.1", "1", "NaN", "invalid"] { + let options = HashMap::from([( + BTREE_INDEX_HIGH_PRIORITY_POOL_RATIO_OPTION.to_string(), + value.to_string(), + )]); + let error = CoreOptions::new(&options) + .btree_index_data_block_cache_size() + .expect_err("invalid BTree cache ratio should fail"); + assert!(error + .to_string() + .contains(BTREE_INDEX_HIGH_PRIORITY_POOL_RATIO_OPTION)); + } + } + #[test] fn test_file_formats_are_normalized_to_lowercase() { // Java routes every file-format option through diff --git a/crates/paimon/src/table/global_index_scanner.rs b/crates/paimon/src/table/global_index_scanner.rs index 8c721b8cd..fdd21a4ea 100644 --- a/crates/paimon/src/table/global_index_scanner.rs +++ b/crates/paimon/src/table/global_index_scanner.rs @@ -40,7 +40,7 @@ use super::global_index_types::{ normalize_queryable_global_index_type, BITMAP_GLOBAL_INDEX_TYPE, BTREE_GLOBAL_INDEX_TYPE, FM_GLOBAL_INDEX_TYPE, MULTIVALUE_GLOBAL_INDEX_TYPE, }; -use crate::btree::{BTreeIndexMeta, BTreeIndexReader}; +use crate::btree::{BTreeDataBlockCache, BTreeIndexMeta, BTreeIndexReader}; use crate::fm_index::{manifest_row_range, FMReadContext, FMReadOptions}; use crate::io::FileIO; use crate::spec::{DataField, FileKind, GlobalIndexSearchMode, IndexManifestEntry, Predicate}; @@ -107,6 +107,7 @@ pub(crate) struct GlobalIndexScanner { /// Scan-scoped shard I/O budget shared by all indexed fields. query_semaphore: Arc, btree_fallback_scan_max_size: i64, + btree_data_block_cache: Arc, bitmap_fallback_scan_max_size: i64, fm_read_options: FMReadOptions, fm_read_context: Arc, @@ -141,6 +142,7 @@ impl GlobalIndexScanner { table_path, global_index_thread_num, btree_fallback_scan_max_size, + 128 * 1024 * 1024, bitmap_fallback_scan_max_size, index_entries, schema_fields, @@ -154,6 +156,7 @@ impl GlobalIndexScanner { table_path: &str, global_index_thread_num: usize, btree_fallback_scan_max_size: i64, + btree_data_block_cache_size: usize, bitmap_fallback_scan_max_size: i64, index_entries: &[IndexManifestEntry], schema_fields: &[DataField], @@ -307,6 +310,7 @@ impl GlobalIndexScanner { global_index_thread_num, query_semaphore: Arc::new(Semaphore::new(global_index_thread_num)), btree_fallback_scan_max_size, + btree_data_block_cache: Arc::new(BTreeDataBlockCache::new(btree_data_block_cache_size)), bitmap_fallback_scan_max_size, fm_read_options, fm_read_context: Arc::new(FMReadContext::new(fm_read_options.cache_size)), @@ -334,6 +338,7 @@ pub(crate) struct GlobalIndexEvaluation<'a> { pub(crate) search_mode: GlobalIndexSearchMode, pub(crate) global_index_thread_num: usize, pub(crate) btree_fallback_scan_max_size: i64, + pub(crate) btree_data_block_cache_size: usize, pub(crate) bitmap_fallback_scan_max_size: i64, pub(crate) fm_read_options: FMReadOptions, pub(crate) next_row_id: Option, @@ -348,6 +353,7 @@ pub(crate) async fn evaluate_global_index( evaluation.table_path, evaluation.global_index_thread_num, evaluation.btree_fallback_scan_max_size, + evaluation.btree_data_block_cache_size, evaluation.bitmap_fallback_scan_max_size, evaluation.index_entries, evaluation.schema_fields, diff --git a/crates/paimon/src/table/global_index_scanner/reader.rs b/crates/paimon/src/table/global_index_scanner/reader.rs index fc011e106..e877176fe 100644 --- a/crates/paimon/src/table/global_index_scanner/reader.rs +++ b/crates/paimon/src/table/global_index_scanner/reader.rs @@ -245,13 +245,20 @@ impl GlobalIndexScanner { let file_reader = input.reader().await?; let cmp = make_key_comparator(data_type); - BTreeIndexReader::open(Box::new(file_reader), file_size, meta, cmp) - .await - .map(OpenedGlobalIndexReader::BTree) - .map_err(|e| crate::Error::DataInvalid { - message: format!("Failed to open BTree index file: {resolved_path}"), - source: Some(Box::new(e)), - }) + BTreeIndexReader::open_with_data_block_cache( + Box::new(file_reader), + file_size, + meta, + cmp, + Arc::clone(&self.btree_data_block_cache), + Arc::from(resolved_path.as_str()), + ) + .await + .map(OpenedGlobalIndexReader::BTree) + .map_err(|e| crate::Error::DataInvalid { + message: format!("Failed to open BTree index file: {resolved_path}"), + source: Some(Box::new(e)), + }) } async fn open_reader_for_entry( diff --git a/crates/paimon/src/table/global_index_scanner/tests.rs b/crates/paimon/src/table/global_index_scanner/tests.rs index fc055c95d..c28a59b32 100644 --- a/crates/paimon/src/table/global_index_scanner/tests.rs +++ b/crates/paimon/src/table/global_index_scanner/tests.rs @@ -693,6 +693,7 @@ async fn evaluate_global_index_fast_with_fallback_size( search_mode: GlobalIndexSearchMode::Fast, global_index_thread_num: 32, btree_fallback_scan_max_size, + btree_data_block_cache_size: 128 * 1024 * 1024, bitmap_fallback_scan_max_size, fm_read_options: FMReadOptions::default(), next_row_id: None, @@ -961,6 +962,7 @@ async fn test_btree_all_match_coverage_nulls_and_boolean_siblings() { search_mode: mode, global_index_thread_num: 2, btree_fallback_scan_max_size: 0, + btree_data_block_cache_size: 128 * 1024 * 1024, bitmap_fallback_scan_max_size: 0, fm_read_options: FMReadOptions::default(), next_row_id: Some(220), @@ -1199,6 +1201,7 @@ async fn test_scalar_optimization_empty_range_never_opens_index() { search_mode: mode, global_index_thread_num: 2, btree_fallback_scan_max_size: 0, + btree_data_block_cache_size: 128 * 1024 * 1024, bitmap_fallback_scan_max_size: 0, fm_read_options: FMReadOptions::default(), next_row_id: Some(220), @@ -2420,6 +2423,7 @@ async fn test_evaluate_global_index_full_mode_includes_unindexed_tail() { search_mode: GlobalIndexSearchMode::Full, global_index_thread_num: 32, btree_fallback_scan_max_size: i64::MAX, + btree_data_block_cache_size: 128 * 1024 * 1024, bitmap_fallback_scan_max_size: i64::MAX, fm_read_options: FMReadOptions::default(), next_row_id: Some(150), @@ -2474,6 +2478,7 @@ async fn test_evaluate_global_index_and_uses_evaluated_field_coverage_for_raw_fa search_mode: GlobalIndexSearchMode::Full, global_index_thread_num: 32, btree_fallback_scan_max_size: i64::MAX, + btree_data_block_cache_size: 128 * 1024 * 1024, bitmap_fallback_scan_max_size: i64::MAX, fm_read_options: FMReadOptions::default(), next_row_id: Some(100), @@ -2509,6 +2514,7 @@ async fn test_evaluate_global_index_detail_mode_uses_data_ranges() { search_mode: GlobalIndexSearchMode::Detail, global_index_thread_num: 32, btree_fallback_scan_max_size: i64::MAX, + btree_data_block_cache_size: 128 * 1024 * 1024, bitmap_fallback_scan_max_size: i64::MAX, fm_read_options: FMReadOptions::default(), next_row_id: Some(150), diff --git a/crates/paimon/src/table/sorted_global_index_build_builder/tests.rs b/crates/paimon/src/table/sorted_global_index_build_builder/tests.rs index aa0d0551b..b19d81ff4 100644 --- a/crates/paimon/src/table/sorted_global_index_build_builder/tests.rs +++ b/crates/paimon/src/table/sorted_global_index_build_builder/tests.rs @@ -793,6 +793,7 @@ async fn assert_btree_build_version(version: u32) { search_mode: GlobalIndexSearchMode::Fast, global_index_thread_num: 32, btree_fallback_scan_max_size: i64::MAX, + btree_data_block_cache_size: 128 * 1024 * 1024, bitmap_fallback_scan_max_size: i64::MAX, fm_read_options: crate::fm_index::FMReadOptions::default(), next_row_id: snapshot.next_row_id(), @@ -1576,6 +1577,7 @@ async fn test_execute_writes_bitmap_index_manifest_and_java_file() { search_mode: GlobalIndexSearchMode::Fast, global_index_thread_num: 32, btree_fallback_scan_max_size: i64::MAX, + btree_data_block_cache_size: 128 * 1024 * 1024, bitmap_fallback_scan_max_size: i64::MAX, fm_read_options: crate::fm_index::FMReadOptions::default(), next_row_id: snapshot.next_row_id(), @@ -1689,6 +1691,7 @@ async fn test_execute_multivalue_index_and_array_queries_end_to_end() { search_mode: GlobalIndexSearchMode::Fast, global_index_thread_num: 32, btree_fallback_scan_max_size: i64::MAX, + btree_data_block_cache_size: 128 * 1024 * 1024, bitmap_fallback_scan_max_size: i64::MAX, fm_read_options: crate::fm_index::FMReadOptions::default(), next_row_id: snapshot.next_row_id(), diff --git a/crates/paimon/src/table/table_scan.rs b/crates/paimon/src/table/table_scan.rs index f4edba05b..c31c51b79 100644 --- a/crates/paimon/src/table/table_scan.rs +++ b/crates/paimon/src/table/table_scan.rs @@ -1765,6 +1765,7 @@ impl<'a> PaimonTableScan<'a> { search_mode: settings.search_mode, global_index_thread_num: settings.thread_num, btree_fallback_scan_max_size: core_options.btree_index_fallback_scan_max_size()?, + btree_data_block_cache_size: core_options.btree_index_data_block_cache_size()?, bitmap_fallback_scan_max_size: core_options .bitmap_index_fallback_scan_max_size()?, fm_read_options: if index_entries