diff --git a/Cargo.lock b/Cargo.lock index ba758429c41fe..e4ec6dae31c7e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5052,6 +5052,7 @@ dependencies = [ "databend-common-sql", "databend-common-statistics", "databend-common-storage", + "databend-common-storages-parquet", "databend-common-tracing", "databend-common-users", "databend-enterprise-fail-safe", @@ -6642,6 +6643,7 @@ dependencies = [ "databend-common-catalog", "databend-common-config", "databend-common-exception", + "databend-common-expression", "databend-common-metrics", "databend-storages-common-index", "databend-storages-common-table-meta", diff --git a/src/query/service/tests/it/storages/fuse/operations/prewhere.rs b/src/query/service/tests/it/storages/fuse/operations/prewhere.rs index 0c089ee730095..1190103e23997 100644 --- a/src/query/service/tests/it/storages/fuse/operations/prewhere.rs +++ b/src/query/service/tests/it/storages/fuse/operations/prewhere.rs @@ -329,6 +329,7 @@ async fn prepare_prewhere_data() -> Result { let part = FuseBlockPartInfo { location: "test_block".to_string(), bloom_filter_index_location: None, + file_size: parquet_bytes.len() as u64, bloom_filter_index_size: 0, create_on: None, nums_rows: num_rows, diff --git a/src/query/storages/common/cache/Cargo.toml b/src/query/storages/common/cache/Cargo.toml index 07799905a83a5..2852d0145edcf 100644 --- a/src/query/storages/common/cache/Cargo.toml +++ b/src/query/storages/common/cache/Cargo.toml @@ -14,6 +14,7 @@ databend-common-cache = { workspace = true } databend-common-catalog = { workspace = true } databend-common-config = { workspace = true } databend-common-exception = { workspace = true } +databend-common-expression = { workspace = true } databend-common-metrics = { workspace = true } databend-storages-common-index = { workspace = true } databend-storages-common-table-meta = { workspace = true } diff --git a/src/query/storages/common/cache/src/caches.rs b/src/query/storages/common/cache/src/caches.rs index 39bf1da621be3..24885216d4090 100644 --- a/src/query/storages/common/cache/src/caches.rs +++ b/src/query/storages/common/cache/src/caches.rs @@ -18,6 +18,7 @@ use std::time::Instant; use arrow::array::ArrayRef; use databend_common_cache::MemSized; +use databend_common_expression::DataBlock; use crate::CacheAccessor; use crate::InMemoryLruCache; @@ -69,10 +70,53 @@ pub type PrunePartitionsCache = InMemoryLruCache<(PartStatistics, Partitions)>; pub type IcebergTableCache = InMemoryLruCache<(Arc, AtomicBool, Instant)>; -/// In memory object cache of table column array -pub type ColumnArrayCache = InMemoryLruCache; -pub type ArrayRawDataUncompressedSize = usize; -pub type SizedColumnArray = (ArrayRef, ArrayRawDataUncompressedSize); +/// In-memory cache of decoded table data. +pub type TableDataCache = InMemoryLruCache; + +pub enum TableDataCacheValue { + ColumnArray(ArrayRef), + DataBlock(DataBlock), +} + +pub struct TableDataCacheEntry { + value: TableDataCacheValue, + memory_size: usize, +} + +impl TableDataCacheEntry { + pub fn from_column_array(value: ArrayRef, memory_size: usize) -> Self { + Self { + value: TableDataCacheValue::ColumnArray(value), + memory_size, + } + } + + pub fn from_data_block(value: DataBlock) -> Self { + let memory_size = value.memory_size(); + Self { + value: TableDataCacheValue::DataBlock(value), + memory_size, + } + } + + pub fn as_column_array(&self) -> Option<&ArrayRef> { + match &self.value { + TableDataCacheValue::ColumnArray(value) => Some(value), + TableDataCacheValue::DataBlock(_) => None, + } + } + + pub fn as_data_block(&self) -> Option<&DataBlock> { + match &self.value { + TableDataCacheValue::ColumnArray(_) => None, + TableDataCacheValue::DataBlock(value) => Some(value), + } + } + + pub fn memory_size(&self) -> usize { + self.memory_size + } +} // Bind Type of cached objects to Caches // @@ -359,10 +403,10 @@ impl From<(PartStatistics, Partitions)> for CacheValue<(PartStatistics, Partitio } } -impl From for CacheValue { - fn from(value: SizedColumnArray) -> Self { +impl From for CacheValue { + fn from(value: TableDataCacheEntry) -> Self { CacheValue { - mem_bytes: value.1, + mem_bytes: value.memory_size, inner: Arc::new(value), } } @@ -384,3 +428,35 @@ impl MemSized for CacheValue { self.mem_bytes } } + +#[cfg(test)] +mod tests { + use arrow::array::Int32Array; + + use super::*; + use crate::CacheAccessor; + + #[test] + fn test_table_data_cache_mixed_entries() { + let cache = TableDataCache::with_bytes_capacity("test_table_data".to_string(), 1024); + + let array: ArrayRef = Arc::new(Int32Array::from(vec![1, 2, 3])); + cache.insert( + "column".to_string(), + TableDataCacheEntry::from_column_array(array.clone(), 128), + ); + let cached_array = cache.get("column").unwrap(); + assert!(cached_array.as_data_block().is_none()); + assert_eq!(cached_array.as_column_array().unwrap().len(), array.len()); + assert_eq!(cached_array.memory_size(), 128); + + cache.insert( + "block".to_string(), + TableDataCacheEntry::from_data_block(DataBlock::empty_with_rows(3)), + ); + let cached_block = cache.get("block").unwrap(); + assert!(cached_block.as_column_array().is_none()); + assert_eq!(cached_block.as_data_block().unwrap().num_rows(), 3); + assert_eq!(cached_block.memory_size(), 0); + } +} diff --git a/src/query/storages/common/cache/src/manager.rs b/src/query/storages/common/cache/src/manager.rs index 39914f449644c..45ec3bfcecbfd 100644 --- a/src/query/storages/common/cache/src/manager.rs +++ b/src/query/storages/common/cache/src/manager.rs @@ -36,7 +36,6 @@ use crate::caches::BlockMetaCache; use crate::caches::BloomIndexFilterCache; use crate::caches::BloomIndexMetaCache; use crate::caches::CacheValue; -use crate::caches::ColumnArrayCache; use crate::caches::ColumnDataCache; use crate::caches::ColumnOrientedSegmentInfoCache; use crate::caches::CompactSegmentInfoCache; @@ -48,6 +47,7 @@ use crate::caches::PrunePartitionsCache; use crate::caches::SegmentBlockMetasCache; use crate::caches::SpatialIndexFileCache; use crate::caches::SpatialIndexMetaCache; +use crate::caches::TableDataCache; use crate::caches::TableSnapshotCache; use crate::caches::TableSnapshotStatisticCache; use crate::caches::VectorIndexFileCache; @@ -118,7 +118,7 @@ pub struct CacheManager { virtual_column_meta_cache: CacheSlot, prune_partitions_cache: CacheSlot, parquet_meta_data_cache: CacheSlot, - in_memory_table_data_cache: CacheSlot, + in_memory_table_data_cache: CacheSlot, segment_block_metas_cache: CacheSlot, block_meta_cache: CacheSlot, @@ -835,7 +835,7 @@ impl CacheManager { self.get_hybrid_cache(self.column_data_cache.get()) } - pub fn get_table_data_array_cache(&self) -> Option { + pub fn get_table_data_cache(&self) -> Option { self.in_memory_table_data_cache.get() } @@ -1023,6 +1023,8 @@ mod tests { use super::*; use crate::ColumnData; + use crate::TableDataCacheEntry; + fn config_with_disk_cache_enabled(cache_path: &str) -> CacheConfig { CacheConfig { data_cache_storage: CacheStorageTypeInnerConfig::Disk, @@ -1259,9 +1261,12 @@ mod tests { // ----- POPULATE BASIC CACHES ----- // 1. Populate in-memory table data cache - let in_memory_table_data_cache = cache_manager.get_table_data_array_cache().unwrap(); + let in_memory_table_data_cache = cache_manager.get_table_data_cache().unwrap(); let array: ArrayRef = Arc::new(Int32Array::from(vec![1, 2, 3, 4, 5])); - in_memory_table_data_cache.insert("not matter".to_string(), (array, 1)); + in_memory_table_data_cache.insert( + "not matter".to_string(), + TableDataCacheEntry::from_column_array(array, 1), + ); // 2. Populate segment block metas cache let segment_block_metas_cache = cache_manager.get_segment_block_metas_cache().unwrap(); @@ -1341,7 +1346,7 @@ mod tests { // ----- VERIFY INITIAL CACHE STATE ----- // Verify all caches are correctly populated - assert!(!cache_manager.get_table_data_array_cache().is_empty()); + assert!(!cache_manager.get_table_data_cache().is_empty()); assert!( !cache_manager .get_segment_block_metas_cache() diff --git a/src/query/storages/fuse/Cargo.toml b/src/query/storages/fuse/Cargo.toml index 0c6d44822b9f0..20f843a059f3f 100644 --- a/src/query/storages/fuse/Cargo.toml +++ b/src/query/storages/fuse/Cargo.toml @@ -25,6 +25,7 @@ databend-common-pipeline-transforms = { workspace = true } databend-common-sql = { workspace = true } databend-common-statistics = { workspace = true } databend-common-storage = { workspace = true } +databend-common-storages-parquet = { workspace = true } databend-common-tracing = { workspace = true } databend-common-users = { workspace = true } databend-enterprise-fail-safe = { workspace = true } diff --git a/src/query/storages/fuse/src/fuse_part.rs b/src/query/storages/fuse/src/fuse_part.rs index e18e306cb6ef0..6dcba37377073 100644 --- a/src/query/storages/fuse/src/fuse_part.rs +++ b/src/query/storages/fuse/src/fuse_part.rs @@ -45,6 +45,8 @@ pub struct FuseBlockPartInfo { pub create_on: Option>, pub nums_rows: usize, + #[serde(default)] + pub file_size: u64, pub columns_meta: HashMap, pub columns_stat: Option>, pub compression: Compression, @@ -96,6 +98,7 @@ impl FuseBlockPartInfo { bloom_filter_index_size, create_on, columns_meta, + file_size: 0, nums_rows: rows_count as usize, compression, sort_min_max, diff --git a/src/query/storages/fuse/src/io/read/block/block_reader.rs b/src/query/storages/fuse/src/io/read/block/block_reader.rs index 01dc040fab256..0a702bbccfbdc 100644 --- a/src/query/storages/fuse/src/io/read/block/block_reader.rs +++ b/src/query/storages/fuse/src/io/read/block/block_reader.rs @@ -272,7 +272,7 @@ impl BlockReadContext { let read_from_in_mem_cache_array: usize = block_read_res .cached_column_array .iter() - .map(|(_, sized_array)| sized_array.1) + .map(|(_, cache_entry)| cache_entry.memory_size()) .sum(); cache_metrics.add_cache_metrics( diff --git a/src/query/storages/fuse/src/io/read/block/block_reader_merge_io.rs b/src/query/storages/fuse/src/io/read/block/block_reader_merge_io.rs index cdb61afcc436a..d8b69991e5edc 100644 --- a/src/query/storages/fuse/src/io/read/block/block_reader_merge_io.rs +++ b/src/query/storages/fuse/src/io/read/block/block_reader_merge_io.rs @@ -19,18 +19,18 @@ use bytes::Bytes; use databend_common_exception::Result; use databend_common_expression::ColumnId; use databend_storages_common_cache::ColumnData; -use databend_storages_common_cache::SizedColumnArray; +use databend_storages_common_cache::TableDataCacheEntry; use databend_storages_common_io::MergeIOReadResult; use enum_as_inner::EnumAsInner; use opendal::Buffer; type CachedColumnData = Vec<(ColumnId, Arc)>; -type CachedColumnArray = Vec<(ColumnId, Arc)>; +type CachedColumnArray = Vec<(ColumnId, Arc)>; #[derive(EnumAsInner, Clone)] pub enum DataItem<'a> { RawData(Buffer), - ColumnArray(&'a Arc), + ColumnArray(&'a Arc), } pub struct BlockReadResult { diff --git a/src/query/storages/fuse/src/io/read/block/block_reader_merge_io_async.rs b/src/query/storages/fuse/src/io/read/block/block_reader_merge_io_async.rs index 1d0be8fa3cb0c..b67ec9aae9016 100644 --- a/src/query/storages/fuse/src/io/read/block/block_reader_merge_io_async.rs +++ b/src/query/storages/fuse/src/io/read/block/block_reader_merge_io_async.rs @@ -64,7 +64,7 @@ impl BlockReadContext { let mut ranges = vec![]; // for async read, try using table data cache (if enabled in settings) let column_data_cache = CacheManager::instance().get_column_data_cache(); - let column_array_cache = CacheManager::instance().get_table_data_array_cache(); + let table_data_cache = CacheManager::instance().get_table_data_cache(); let mut cached_column_data = vec![]; let mut cached_column_array = vec![]; @@ -84,14 +84,16 @@ impl BlockReadContext { // first, check in memory table data cache // column_array_cache - if let Some(cache_array) = column_array_cache.get_sized(&column_cache_key, len) { - // Record bytes scanned from memory cache (table data only) - Profile::record_usize_profile( - ProfileStatisticsName::ScanBytesFromMemory, - len as usize, - ); - cached_column_array.push((*column_id, cache_array)); - continue; + if let Some(cache_entry) = table_data_cache.get_sized(&column_cache_key, len) { + if cache_entry.as_column_array().is_some() { + // Record bytes scanned from memory cache (table data only) + Profile::record_usize_profile( + ProfileStatisticsName::ScanBytesFromMemory, + len as usize, + ); + cached_column_array.push((*column_id, cache_entry)); + continue; + } } // and then, check on disk table data cache diff --git a/src/query/storages/fuse/src/io/read/block/parquet/mod.rs b/src/query/storages/fuse/src/io/read/block/parquet/mod.rs index aaeb22d5291b3..4d69172714515 100644 --- a/src/query/storages/fuse/src/io/read/block/parquet/mod.rs +++ b/src/query/storages/fuse/src/io/read/block/parquet/mod.rs @@ -13,11 +13,14 @@ // limitations under the License. use std::collections::HashMap; +use std::fmt::Write; use arrow_array::Array; use arrow_array::ArrayRef; use arrow_array::RecordBatch; use arrow_array::StructArray; +use databend_common_base::runtime::profile::Profile; +use databend_common_base::runtime::profile::ProfileStatisticsName; use databend_common_catalog::plan::Projection; use databend_common_exception::ErrorCode; use databend_common_expression::BlockEntry; @@ -27,9 +30,11 @@ use databend_common_expression::FilterVisitor; use databend_common_expression::TableDataType; use databend_common_expression::TableSchema; use databend_common_expression::Value; +use databend_common_expression::types::DataType; use databend_common_expression::visitor::ValueVisitor; use databend_storages_common_cache::CacheAccessor; use databend_storages_common_cache::CacheManager; +use databend_storages_common_cache::TableDataCacheEntry; use databend_storages_common_cache::TableDataCacheKey; use databend_storages_common_table_meta::meta::ColumnMeta; use databend_storages_common_table_meta::meta::Compression; @@ -62,6 +67,124 @@ impl BlockReader { ) } + pub fn page_range_bitmap( + part: &FuseBlockPartInfo, + ) -> Option { + part.range().map(|range| { + RowSelection::from_range( + part.nums_rows, + range.start.saturating_mul(part.page_size()), + range.end.saturating_mul(part.page_size()), + ) + .bitmap + }) + } + + pub fn page_range_data_cache_key(&self, part: &FuseBlockPartInfo) -> Option { + let range = part.range()?; + let root = self.operator.info().root(); + let mut key = String::with_capacity(root.len() + part.location.len() + 256); + write!( + key, + "fuse-page-data-v1|{}:{}|{}:{}|{}|{}|{}|{}:{}", + root.len(), + root, + part.location.len(), + part.location, + part.file_size, + part.nums_rows, + part.page_size(), + range.start, + range.end, + ) + .ok()?; + + for (field, column_node) in self + .projected_schema + .fields + .iter() + .zip(self.project_column_nodes.iter()) + { + if column_node.is_nested || column_node.leaf_column_ids.as_slice() != [field.column_id] + { + return None; + } + let column_meta = part.columns_meta.get(&field.column_id)?; + let (offset, len) = column_meta.offset_length(); + let data_type: DataType = field.data_type().into(); + write!( + key, + "|{}:{}:{}:{:?}", + field.column_id, offset, len, data_type + ) + .ok()?; + } + Some(key) + } + + pub fn cached_page_range_data(&self, key: &str) -> Option { + let cache = CacheManager::instance().get_table_data_cache(); + let cache_entry = cache.get(key)?; + let data_block = cache_entry.as_data_block()?; + let memory_size = cache_entry.memory_size(); + Profile::record_usize_profile(ProfileStatisticsName::ScanBytesFromMemory, memory_size); + self.ctx + .get_data_cache_metrics() + .add_cache_metrics(0, 0, memory_size); + Some(data_block.clone()) + } + + pub fn cache_page_range_data(&self, key: String, data_block: DataBlock) -> DataBlock { + if !self.put_cache { + return data_block; + } + let Some(cache) = CacheManager::instance().get_table_data_cache() else { + return data_block; + }; + let cache_entry = cache.insert(key, TableDataCacheEntry::from_data_block(data_block)); + cache_entry.as_data_block().unwrap().clone() + } + + pub fn deserialize_parquet_record_batch( + &self, + part: &FuseBlockPartInfo, + record_batch: &RecordBatch, + ) -> databend_common_exception::Result { + let result_rows = record_batch.num_rows(); + if self.projected_schema.fields.is_empty() { + return Ok(DataBlock::empty_with_rows(result_rows)); + } + + if result_rows == 0 { + return Ok(DataBlock::empty_with_schema(&self.data_schema())); + } + + let mut entries = Vec::with_capacity(self.projected_schema.fields.len()); + let name_paths = column_name_paths(&self.projection, &self.original_schema); + for (((i, field), column_node), name_path) in self + .projected_schema + .fields + .iter() + .enumerate() + .zip(self.project_column_nodes.iter()) + .zip(name_paths.iter()) + { + let data_type = field.data_type().into(); + let exists = column_node + .leaf_column_ids + .iter() + .any(|column_id| part.columns_meta.contains_key(column_id)); + let value = if exists { + let arrow_array = column_by_name(record_batch, name_path); + Value::from_arrow_rs(arrow_array, &data_type)? + } else { + Value::Scalar(self.default_vals[i].clone()) + }; + entries.push(BlockEntry::new(value, || (data_type, result_rows))); + } + Ok(DataBlock::new(entries, result_rows)) + } + pub fn deserialize_parquet_chunks( &self, num_rows: usize, @@ -94,7 +217,7 @@ impl BlockReader { let name_paths = column_name_paths(&self.projection, &self.original_schema); let array_cache = if self.put_cache && !has_selection { - CacheManager::instance().get_table_data_array_cache() + CacheManager::instance().get_table_data_cache() } else { None }; @@ -133,7 +256,13 @@ impl BlockReader { let key = TableDataCacheKey::new(block_path, field.column_id, offset, len); let array_memory_size = arrow_array.get_array_memory_size(); - cache.insert(key.into(), (arrow_array.clone(), array_memory_size)); + cache.insert( + key.into(), + TableDataCacheEntry::from_column_array( + arrow_array.clone(), + array_memory_size, + ), + ); } } Value::from_arrow_rs(arrow_array, &data_type)? @@ -145,7 +274,10 @@ impl BlockReader { "unexpected nested field: nested leaf field hits cached", )); } - let mut value = Value::from_arrow_rs(cached.0.clone(), &data_type)?; + let cached_array = cached.as_column_array().ok_or_else(|| { + ErrorCode::StorageOther("unexpected data block entry in column array cache") + })?; + let mut value = Value::from_arrow_rs(cached_array.clone(), &data_type)?; if let Some(selection) = selection { let mut filter_visitor = FilterVisitor::new(&selection.bitmap); filter_visitor.visit_value(value)?; diff --git a/src/query/storages/fuse/src/io/read/block/parquet/row_selection.rs b/src/query/storages/fuse/src/io/read/block/parquet/row_selection.rs index 4981db3d9fd59..b58289ef13390 100644 --- a/src/query/storages/fuse/src/io/read/block/parquet/row_selection.rs +++ b/src/query/storages/fuse/src/io/read/block/parquet/row_selection.rs @@ -14,6 +14,7 @@ use databend_common_column::bitmap::utils::SlicesIterator; use databend_common_expression::types::Bitmap; +use databend_common_expression::types::MutableBitmap; use parquet::arrow::arrow_reader::RowSelector; /// A wrapper around parquet's `RowSelection` that also tracks the number of selected rows (bits set to 1 in the bitmap). @@ -38,6 +39,37 @@ impl RowSelection { bitmap, } } + + pub fn from_range(total_rows: usize, start: usize, end: usize) -> Self { + let start = start.min(total_rows); + let end = end.min(total_rows).max(start); + + let selected_rows = end - start; + let mut bitmap = MutableBitmap::with_capacity(total_rows); + bitmap.extend_constant(start, false); + bitmap.extend_constant(selected_rows, true); + bitmap.extend_constant(total_rows - end, false); + + let mut selectors = Vec::with_capacity(3); + if selected_rows == 0 { + if total_rows > 0 { + selectors.push(RowSelector::skip(total_rows)); + } + } else { + if start > 0 { + selectors.push(RowSelector::skip(start)); + } + selectors.push(RowSelector::select(selected_rows)); + if end < total_rows { + selectors.push(RowSelector::skip(total_rows - end)); + } + } + Self::new( + parquet::arrow::arrow_reader::RowSelection::from(selectors), + selected_rows, + bitmap.into(), + ) + } } impl From<&Bitmap> for RowSelection { @@ -136,4 +168,25 @@ mod tests { assert_eq!(selectors[0].row_count, 5); assert!(selectors[0].skip); } + + #[test] + fn test_row_selection_from_range() { + let row_selection = RowSelection::from_range(10, 2, 5); + let selectors: Vec<_> = row_selection.selection.iter().collect(); + + assert_eq!(row_selection.selected_rows, 3); + assert_eq!(row_selection.bitmap.iter().collect::>(), vec![ + false, false, true, true, true, false, false, false, false, false + ]); + assert_eq!(selectors.len(), 3); + assert_eq!((selectors[0].row_count, selectors[0].skip), (2, true)); + assert_eq!((selectors[1].row_count, selectors[1].skip), (3, false)); + assert_eq!((selectors[2].row_count, selectors[2].skip), (5, true)); + + let empty = RowSelection::from_range(10, 4, 4); + let selectors: Vec<_> = empty.selection.iter().collect(); + assert_eq!(empty.selected_rows, 0); + assert_eq!(selectors.len(), 1); + assert_eq!((selectors[0].row_count, selectors[0].skip), (10, true)); + } } diff --git a/src/query/storages/fuse/src/operations/read/fuse_rows_fetcher.rs b/src/query/storages/fuse/src/operations/read/fuse_rows_fetcher.rs index 98b4000a8c210..64020811423d6 100644 --- a/src/query/storages/fuse/src/operations/read/fuse_rows_fetcher.rs +++ b/src/query/storages/fuse/src/operations/read/fuse_rows_fetcher.rs @@ -73,6 +73,7 @@ pub fn row_fetch_processor( FuseStorageFormat::Parquet => { let read_settings = ReadSettings::from_ctx(&ctx)?; let max_threads = ctx.get_settings().get_max_threads()? as usize; + let source_parts = source.parts.partitions.clone(); let block_threshold = BlockThreshold { max_rows: ctx.get_settings().get_max_block_size()? as usize, max_bytes: ctx.get_settings().get_max_block_bytes()? as usize, @@ -92,6 +93,7 @@ pub fn row_fetch_processor( read_settings, max_threads, io_semaphore.clone(), + &source_parts, ), need_wrap_nullable, fetched_data_types.clone(), diff --git a/src/query/storages/fuse/src/operations/read/parquet_rows_fetcher.rs b/src/query/storages/fuse/src/operations/read/parquet_rows_fetcher.rs index 375acb536a130..bfdc39999a09f 100644 --- a/src/query/storages/fuse/src/operations/read/parquet_rows_fetcher.rs +++ b/src/query/storages/fuse/src/operations/read/parquet_rows_fetcher.rs @@ -14,10 +14,14 @@ use std::collections::HashMap; use std::collections::hash_map::Entry; +use std::fmt::Write; use std::future::Future; +use std::ops::Range; use std::sync::Arc; +use std::sync::LazyLock; use databend_common_base::runtime::spawn; +use databend_common_catalog::plan::PartInfoPtr; use databend_common_catalog::plan::Projection; use databend_common_catalog::plan::block_id_in_segment; use databend_common_catalog::plan::block_idx_in_segment; @@ -36,6 +40,7 @@ use databend_storages_common_cache::CacheValue; use databend_storages_common_cache::InMemoryLruCache; use databend_storages_common_cache::LoadParams; use databend_storages_common_io::ReadSettings; +use databend_storages_common_pruner::BlockMetaIndex; use databend_storages_common_table_meta::meta::BlockMeta; use databend_storages_common_table_meta::meta::ColumnMeta; use databend_storages_common_table_meta::meta::Compression; @@ -48,7 +53,9 @@ use tokio::task::JoinHandle; use super::fuse_rows_fetcher::RowsFetchMetadata; use super::fuse_rows_fetcher::RowsFetcher; +use super::read_block_context::read_parquet_page_range_data; use crate::BlockReadResult; +use crate::FUSE_OPT_KEY_DATA_PAGE_ROWS; use crate::FuseBlockPartInfo; use crate::FuseTable; use crate::io::BlockReader; @@ -95,10 +102,20 @@ pub struct RowsFetchMetadataImpl { pub location: String, pub nums_rows: usize, + pub file_size: u64, pub compression: Compression, pub columns_meta: HashMap, + pub block_meta_index: Option, } +static ROW_FETCH_METADATA_CACHE: LazyLock> = + LazyLock::new(|| { + InMemoryLruCache::with_items_capacity( + String::from("RowFetchProjectedBlockMetaCache"), + 16_384, + ) + }); + impl RowsFetchMetadata for RowsFetchMetadataImpl { fn row_bytes(&self) -> usize { self.row_bytes @@ -127,6 +144,8 @@ pub(super) struct ParquetRowsFetcher { segment_reader: CompactSegmentInfoReader, block_meta_lru_cache: InMemoryLruCache, + source_block_meta_indexes: HashMap, + metadata_cache_prefix: String, } #[async_trait::async_trait] @@ -135,16 +154,22 @@ impl RowsFetcher for ParquetRowsFetcher { #[async_backtrace::framed] async fn initialize(&mut self) -> Result<()> { - self.snapshot = self.table.read_table_snapshot().await?; Ok(()) } async fn fetch_metadata(&mut self, block_id: u64) -> Result { + let global_key = self.metadata_cache_key(block_id); + if let Some(metadata) = ROW_FETCH_METADATA_CACHE.get(&global_key) { + return Ok(metadata); + } if let Some(v) = self.block_meta_lru_cache.get(block_id.to_string()) { return Ok(v.clone()); } // load metadata let (segment, block) = split_prefix(block_id); + if self.snapshot.is_none() { + self.snapshot = self.table.read_table_snapshot().await?; + } let snapshot = self.snapshot.as_ref().unwrap(); let (location, ver) = snapshot.segments[segment as usize].clone(); @@ -166,18 +191,22 @@ impl RowsFetcher for ParquetRowsFetcher { cache_end = std::cmp::min(cache_end, blocks.len()); let metadata = self.build_metadata(&blocks[cache_start..cache_end])?; - for (block_index, metadata) in (cache_start..cache_end).zip(metadata.into_iter()) { + for (block_index, mut metadata) in (cache_start..cache_end).zip(metadata.into_iter()) { let block_id = block_id_in_segment(blocks.len(), block_index); let block_id = compute_row_id_prefix(segment, block_id as u64); - self.block_meta_lru_cache - .insert(block_id.to_string(), metadata); + metadata.block_meta_index = self.source_block_meta_indexes.get(&block_id).cloned(); + if metadata.block_meta_index.is_some() { + ROW_FETCH_METADATA_CACHE.insert(self.metadata_cache_key(block_id), metadata); + } else { + self.block_meta_lru_cache + .insert(block_id.to_string(), metadata); + } } - Ok(self - .block_meta_lru_cache - .get(block_id.to_string()) - .clone() - .unwrap()) + ROW_FETCH_METADATA_CACHE + .get(global_key) + .or_else(|| self.block_meta_lru_cache.get(block_id.to_string())) + .ok_or_else(|| ErrorCode::Internal("Row fetch block metadata is missing")) } #[async_backtrace::framed] @@ -288,10 +317,25 @@ impl ParquetRowsFetcher { settings: ReadSettings, max_threads: usize, io_semaphore: Arc, + source_parts: &[PartInfoPtr], ) -> Self { let schema = table.schema(); let operator = table.operator.clone(); let segment_reader = MetaReaders::segment_info_reader(operator, schema.clone()); + let metadata_cache_prefix = projection_cache_prefix(&table, &projection); + let data_page_rows = table.get_option(FUSE_OPT_KEY_DATA_PAGE_ROWS, 0usize); + let source_block_meta_indexes = source_parts + .iter() + .filter_map(|part| FuseBlockPartInfo::from_part(part).ok()) + .filter_map(|part| part.block_meta_index.clone()) + .map(|mut index| { + if data_page_rows > 0 { + index.page_size = data_page_rows; + } + let prefix = compute_row_id_prefix(index.segment_idx as u64, index.block_id as u64); + (prefix, index) + }) + .collect(); ParquetRowsFetcher { table, snapshot: None, @@ -306,9 +350,15 @@ impl ParquetRowsFetcher { String::from("RowFetchBlockMetaCache"), 128, ), + source_block_meta_indexes, + metadata_cache_prefix, } } + fn metadata_cache_key(&self, block_id: u64) -> String { + format!("{}{}", self.metadata_cache_prefix, block_id) + } + fn fetch_block( &self, metadata: Arc, @@ -327,6 +377,13 @@ impl ParquetRowsFetcher { .acquire_owned() .await .expect("row-fetch io semaphore never closed"); + if let Some(block) = + Self::fetch_page_range(&reader, &settings, &metadata, take_indices.as_slice()) + .await? + { + return Ok((final_index, block)); + } + let chunk = reader .read_columns_data_by_merge_io( &settings, @@ -344,6 +401,47 @@ impl ParquetRowsFetcher { } } + async fn fetch_page_range( + reader: &Arc, + settings: &ReadSettings, + metadata: &RowsFetchMetadataImpl, + take_indices: &[u32], + ) -> Result> { + let Some(mut block_meta_index) = metadata.block_meta_index.clone() else { + return Ok(None); + }; + let Some((range, range_start, relative_indices)) = + page_range_for_rows(take_indices, metadata.nums_rows, block_meta_index.page_size)? + else { + return Ok(None); + }; + + block_meta_index.range = Some(range); + let part = FuseBlockPartInfo { + location: metadata.location.clone(), + bloom_filter_index_location: None, + bloom_filter_index_size: 0, + create_on: None, + nums_rows: metadata.nums_rows, + file_size: metadata.file_size, + columns_meta: metadata.columns_meta.clone(), + columns_stat: None, + compression: metadata.compression, + sort_min_max: None, + block_meta_index: Some(block_meta_index), + }; + + let Some(block) = read_parquet_page_range_data(reader, settings, &part).await? else { + return Ok(None); + }; + debug_assert!( + relative_indices + .iter() + .all(|index| *index as usize + range_start < metadata.nums_rows) + ); + Ok(Some(block.take(relative_indices.as_slice())?)) + } + fn build_metadata(&self, meta: &[Arc]) -> Result> { let arrow_schema = self.schema.as_ref().into(); let column_nodes = ColumnNodes::new_from_schema(&arrow_schema, Some(&self.schema)); @@ -383,9 +481,11 @@ impl ParquetRowsFetcher { row_bytes: average_bytes, block_bytes, nums_rows: fuse_part.nums_rows, + file_size: block_meta.file_size, compression: fuse_part.compression, location: fuse_part.location.clone(), columns_meta: fuse_part.columns_meta.clone(), + block_meta_index: None, }); } @@ -409,6 +509,86 @@ impl ParquetRowsFetcher { } } +fn projection_cache_prefix(table: &FuseTable, projection: &Projection) -> String { + let table_ident = table.get_table_info().ident; + metadata_cache_prefix( + &table.get_operator_ref().info().root(), + table_ident.table_id, + table_ident.seq, + &table.query_result_cache_id(), + table.get_option(FUSE_OPT_KEY_DATA_PAGE_ROWS, 0usize), + projection, + ) +} + +fn metadata_cache_prefix( + operator_root: &str, + table_id: u64, + table_seq: u64, + snapshot_cache_id: &str, + data_page_rows: usize, + projection: &Projection, +) -> String { + let mut key = format!( + "row-fetch-meta-v3|{}:{}|{}|{}|{}:{}|{}|", + operator_root.len(), + operator_root, + table_id, + table_seq, + snapshot_cache_id.len(), + snapshot_cache_id, + data_page_rows, + ); + match projection { + Projection::Columns(indices) => { + key.push('c'); + for index in indices { + write!(key, ":{index}").unwrap(); + } + } + Projection::InnerColumns(paths) => { + key.push('i'); + for (index, path) in paths { + write!(key, ":{index}=").unwrap(); + for child in path { + write!(key, "{child},").unwrap(); + } + } + } + } + key.push(':'); + key +} + +type PageRangeSelection = (Range, usize, Vec); + +fn page_range_for_rows( + take_indices: &[u32], + nums_rows: usize, + page_size: usize, +) -> Result> { + if take_indices.is_empty() || page_size == 0 || page_size >= nums_rows { + return Ok(None); + } + + let min_row = *take_indices.iter().min().unwrap() as usize; + let max_row = *take_indices.iter().max().unwrap() as usize; + if max_row >= nums_rows { + return Err(ErrorCode::Internal(format!( + "RowID index {max_row} is outside a block with {nums_rows} rows" + ))); + } + + let start_page = min_row / page_size; + let end_page = max_row / page_size + 1; + let range_start = start_page * page_size; + let relative_indices = take_indices + .iter() + .map(|index| *index - range_start as u32) + .collect(); + Ok(Some((start_page..end_page, range_start, relative_indices))) +} + impl From for CacheValue { fn from(value: RowsFetchMetadataImpl) -> Self { CacheValue::new(value, 0) @@ -481,4 +661,59 @@ mod tests { drop(permit); assert_eq!(semaphore.available_permits(), 1); } + + #[test] + fn test_page_range_for_rows_preserves_output_order() { + let (range, start, relative) = page_range_for_rows(&[4100, 4097, 8192, 4100], 12_000, 4096) + .unwrap() + .unwrap(); + + assert_eq!(range, 1..3); + assert_eq!(start, 4096); + assert_eq!(relative, vec![4, 1, 4096, 4]); + } + + #[test] + fn test_page_range_for_rows_rejects_invalid_row_id() { + let err = page_range_for_rows(&[10], 10, 4).unwrap_err(); + assert!(err.message().contains("outside a block")); + } + + #[test] + fn test_metadata_cache_prefix_isolates_storage_and_snapshot() { + let projection = Projection::Columns(vec![1, 3]); + let base = metadata_cache_prefix("root-a", 1, 2, "snapshot-a", 100, &projection); + + assert_ne!( + base, + metadata_cache_prefix("root-b", 1, 2, "snapshot-a", 100, &projection) + ); + assert_ne!( + base, + metadata_cache_prefix("root-a", 1, 2, "snapshot-b", 100, &projection) + ); + assert_ne!( + base, + metadata_cache_prefix("root-a", 2, 2, "snapshot-a", 100, &projection) + ); + assert_ne!( + base, + metadata_cache_prefix("root-a", 1, 3, "snapshot-a", 100, &projection) + ); + assert_ne!( + base, + metadata_cache_prefix("root-a", 1, 2, "snapshot-a", 200, &projection) + ); + assert_ne!( + metadata_cache_prefix("ab", 1, 2, "c", 100, &projection), + metadata_cache_prefix("a", 1, 2, "bc", 100, &projection) + ); + } + + #[test] + fn test_page_range_for_rows_skips_empty_or_unknown_pages() { + assert!(page_range_for_rows(&[], 10, 4).unwrap().is_none()); + assert!(page_range_for_rows(&[1], 10, 0).unwrap().is_none()); + assert!(page_range_for_rows(&[1], 10, 10).unwrap().is_none()); + } } diff --git a/src/query/storages/fuse/src/operations/read/read_block_context.rs b/src/query/storages/fuse/src/operations/read/read_block_context.rs index a0121ee638a9e..c805575f30c7c 100644 --- a/src/query/storages/fuse/src/operations/read/read_block_context.rs +++ b/src/query/storages/fuse/src/operations/read/read_block_context.rs @@ -14,11 +14,25 @@ use std::sync::Arc; +use arrow_schema::Schema; use databend_common_catalog::plan::PartInfoPtr; use databend_common_catalog::table_context::TableContext; +use databend_common_exception::ErrorCode; use databend_common_exception::Result; +use databend_common_expression::DataBlock; +use databend_common_storages_parquet::InMemoryRowGroup; +use databend_common_storages_parquet::ParquetFileReader; +use databend_common_storages_parquet::ReadSettings as ParquetReadSettings; +use databend_storages_common_cache::CacheAccessor; +use databend_storages_common_cache::CacheManager; use databend_storages_common_io::ReadSettings; use log::debug; +use parquet::arrow::ProjectionMask; +use parquet::arrow::arrow_reader::ParquetRecordBatchReader; +use parquet::arrow::parquet_to_arrow_field_levels; +use parquet::file::metadata::PageIndexPolicy; +use parquet::file::metadata::ParquetMetaData; +use parquet::file::metadata::ParquetMetaDataReader; use super::block_format::FuseParquetBlockFormat; use super::parquet_data_source::ParquetDataSource; @@ -26,6 +40,8 @@ use crate::FuseBlockPartInfo; use crate::FuseStorageFormat; use crate::io::AggIndexReader; use crate::io::BlockReadContext; +use crate::io::BlockReader; +use crate::io::RowSelection; use crate::io::TableMetaLocationGenerator; use crate::io::VirtualBlockReadResult; use crate::io::VirtualColumnReader; @@ -168,3 +184,117 @@ impl ReadBlockContext { .await } } + +pub(crate) async fn read_parquet_page_range_data( + block_reader: &Arc, + read_settings: &ReadSettings, + part: &FuseBlockPartInfo, +) -> Result> { + let Some(page_cache_key) = block_reader.page_range_data_cache_key(part) else { + return Ok(None); + }; + if let Some(data_block) = block_reader.cached_page_range_data(&page_cache_key) { + return Ok(Some(data_block)); + } + + let block_read_ctx = block_reader.read_context(); + let metadata = parquet_metadata_with_offset_indexes(&block_read_ctx, part).await?; + if metadata.num_row_groups() != 1 { + return Ok(None); + } + let Some(offset_indexes) = metadata.offset_index().and_then(|v| v.first()) else { + return Ok(None); + }; + + let schema_descr = metadata.file_metadata().schema_descr(); + let projection_indices = block_read_ctx + .project_indices() + .iter() + .filter_map(|(index, (column_id, ..))| { + part.columns_meta.contains_key(column_id).then_some(*index) + }) + .collect::>(); + if projection_indices.is_empty() + || projection_indices + .iter() + .any(|index| *index >= schema_descr.num_columns()) + { + return Ok(None); + } + + let page_bitmap = BlockReader::page_range_bitmap(part) + .ok_or_else(|| ErrorCode::Internal("page range is missing"))?; + let parquet_selection = RowSelection::from(&page_bitmap).selection; + let page_locations = offset_indexes + .iter() + .map(|index| index.page_locations().to_vec()) + .collect::>(); + if page_locations.len() != schema_descr.num_columns() { + return Ok(None); + } + + let projection = ProjectionMask::leaves(schema_descr, projection_indices); + let parquet_read_settings = ParquetReadSettings { + max_gap_size: read_settings.max_gap_size, + max_range_size: read_settings.max_range_size, + parquet_fast_read_bytes: read_settings.parquet_fast_read_bytes, + enable_cache: true, + }; + let mut row_group = InMemoryRowGroup::new( + &part.location, + block_read_ctx.operator().clone(), + metadata.row_group(0), + Some(page_locations), + parquet_read_settings, + ); + row_group + .fetch(&projection, Some(&parquet_selection)) + .await?; + + let arrow_schema = Schema::from(block_reader.original_schema.as_ref()); + let field_levels = + parquet_to_arrow_field_levels(schema_descr, projection, Some(arrow_schema.fields()))?; + let mut reader = ParquetRecordBatchReader::try_new_with_row_groups( + &field_levels, + &row_group, + part.nums_rows, + Some(parquet_selection), + )?; + let record_batch = reader + .next() + .ok_or_else(|| ErrorCode::Internal("selected parquet range returned no rows"))??; + debug_assert!(reader.next().is_none()); + + let data_block = block_reader.deserialize_parquet_record_batch(part, &record_batch)?; + Ok(Some( + block_reader.cache_page_range_data(page_cache_key, data_block), + )) +} + +async fn parquet_metadata_with_offset_indexes( + block_read_ctx: &BlockReadContext, + part: &FuseBlockPartInfo, +) -> Result> { + let cache = CacheManager::instance().get_parquet_meta_data_cache(); + let cache_key = format!( + "{}{}", + block_read_ctx.operator().info().root(), + part.location + ); + if let Some(metadata) = cache.as_ref().and_then(|cache| cache.get(&cache_key)) { + if metadata.offset_index().is_some() { + return Ok(metadata); + } + } + + let op_reader = block_read_ctx.operator().reader(&part.location).await?; + let mut file_reader = ParquetFileReader::new(op_reader, part.file_size); + let metadata = ParquetMetaDataReader::new() + .with_offset_index_policy(PageIndexPolicy::Optional) + .load_and_finish(&mut file_reader, part.file_size) + .await?; + Ok(match cache { + Some(cache) => cache.insert(cache_key, metadata), + None => Arc::new(metadata), + }) +} diff --git a/src/query/storages/system/src/caches_table.rs b/src/query/storages/system/src/caches_table.rs index d4245244ca728..8bdc5ea652978 100644 --- a/src/query/storages/system/src/caches_table.rs +++ b/src/query/storages/system/src/caches_table.rs @@ -93,7 +93,7 @@ impl SyncSystemTable for CachesTable { let prune_partitions_cache = cache_manager.get_prune_partitions_cache(); let parquet_meta_data_cache = cache_manager.get_parquet_meta_data_cache(); let column_data_cache = cache_manager.get_column_data_cache(); - let table_column_array_cache = cache_manager.get_table_data_array_cache(); + let table_data_cache = cache_manager.get_table_data_cache(); let iceberg_table_cache = cache_manager.get_iceberg_table_cache(); let mut columns = CachesTableColumns::default(); @@ -189,8 +189,8 @@ impl SyncSystemTable for CachesTable { Self::append_row(&parquet_meta_data_cache, &local_node, &mut columns); } - if let Some(table_column_array_cache) = table_column_array_cache { - Self::append_row(&table_column_array_cache, &local_node, &mut columns); + if let Some(table_data_cache) = table_data_cache { + Self::append_row(&table_data_cache, &local_node, &mut columns); } if let Some(iceberg_table_cache) = iceberg_table_cache { diff --git a/tests/task/test-private-task.sh b/tests/task/test-private-task.sh index 92ab0a19ad054..6ab4dd6a4c347 100644 --- a/tests/task/test-private-task.sh +++ b/tests/task/test-private-task.sh @@ -706,6 +706,30 @@ else exit 1 fi +# Ensure the rerun parent completes in a later second than every previous child run. +actual_child_success_count=0 +rerun_clock_ready=0 +for _ in {1..20}; do + response=$(query_sql_with_auth "root:" "SELECT count(*), if(to_start_of_second(now()) > max(completed_at), 1, 0) FROM system_task.task_run WHERE task_name IN ('fanout_child_a', 'fanout_child_b', 'fanout_child_c') AND state = 'SUCCEEDED' AND completed_at IS NOT NULL") + check_response_error "$response" + actual_child_success_count=$(echo "$response" | jq -r '.data[0][0]') + rerun_clock_ready=$(echo "$response" | jq -r '.data[0][1]') + if [ "$actual_child_success_count" = "3" ] && [ "$rerun_clock_ready" = "1" ]; then + break + fi + sleep 1 +done + +if [ "$actual_child_success_count" = "3" ] && [ "$rerun_clock_ready" = "1" ]; then + echo "✅ Private task fan-out child runs are terminal and rerun clock is ready" +else + echo "❌ Expected terminal private task fan-out children and a later rerun clock" + echo "Expected child successes: 3" + echo "Actual child successes : $actual_child_success_count" + echo "Rerun clock ready : $rerun_clock_ready" + exit 1 +fi + response=$(query_sql_with_auth "root:" "UPDATE system_task.task_run SET state = 'SKIPPED', error_code = 0, error_message = 'OVERLAPPING_EXECUTION: test', completed_at = to_timestamp(4102444800) WHERE task_name = 'fanout_root'") check_response_error "$response"