From 26c3281f2ba71e33be6a02e79d5f5499371bdbb0 Mon Sep 17 00:00:00 2001 From: Yingwen Date: Thu, 17 Sep 2026 08:45:00 +0000 Subject: [PATCH] perf(mito2): use series and range indexes in two-phase series scans (#9153) * feat: use series indexes for SeriesScan candidate discovery Signed-off-by: evenyag * perf: use range indexes in two-phase series reader Signed-off-by: evenyag * fix(mito2): update series index test fixtures after rebase Signed-off-by: evenyag * refactor(mito2): move lazy range index searchers into file contexts Signed-off-by: evenyag * fix: pin series index snapshot before data snapshot Signed-off-by: evenyag * fix: preserve builder caching for index-covered SSTs Signed-off-by: evenyag --------- Signed-off-by: evenyag --- src/mito2/src/engine.rs | 10 + src/mito2/src/engine/scan_test.rs | 172 ++++++++- src/mito2/src/read/pruner.rs | 8 +- src/mito2/src/read/scan_region.rs | 29 ++ src/mito2/src/read/series_candidate.rs | 430 ++++++++++++++++++++- src/mito2/src/read/series_reader.rs | 34 +- src/mito2/src/region.rs | 5 + src/mito2/src/region/opener.rs | 2 + src/mito2/src/series_index.rs | 20 +- src/mito2/src/series_index/catalog.rs | 61 +++ src/mito2/src/series_index/searcher.rs | 172 ++++++--- src/mito2/src/series_index/version.rs | 5 + src/mito2/src/sst/parquet/file_range.rs | 154 +++++++- src/mito2/src/sst/parquet/reader.rs | 25 ++ src/mito2/src/sst/parquet/row_selection.rs | 2 +- 15 files changed, 1038 insertions(+), 91 deletions(-) diff --git a/src/mito2/src/engine.rs b/src/mito2/src/engine.rs index 986e205ee39..e1458297a16 100644 --- a/src/mito2/src/engine.rs +++ b/src/mito2/src/engine.rs @@ -1053,6 +1053,15 @@ impl EngineInner { let query_start = Instant::now(); // Reading a region doesn't need to go through the region worker thread. let region = self.find_region(region_id)?; + // Pin the index before the data snapshot: compaction and index publication + // could otherwise give us a newer index that omits series still visible + // in the query's older SST snapshot. + let series_index = region.series_index_store.as_ref().map(|store| { + crate::series_index::SeriesIndexReadContext { + store: store.clone(), + version: region.series_index_version(), + } + }); let version_data = region.version_control.current(); let version = version_data.version; @@ -1146,6 +1155,7 @@ impl EngineInner { request, CacheStrategy::EnableAll(cache_manager), ) + .with_series_index(series_index) .with_query_stat_counters(region.region_stats.query_stat_counters()) .with_max_concurrent_scan_files(self.config.max_concurrent_scan_files) .with_scan_memory_pool(self.scan_memory_pool.clone()) diff --git a/src/mito2/src/engine/scan_test.rs b/src/mito2/src/engine/scan_test.rs index 0dc20d97107..e5362aa4ddf 100644 --- a/src/mito2/src/engine/scan_test.rs +++ b/src/mito2/src/engine/scan_test.rs @@ -1232,10 +1232,19 @@ async fn test_series_scan_with_format(flat_format: bool) { #[tokio::test] async fn test_two_phase_series_scan() { + for use_index in [false, true] { + for use_range_index in [false, true] { + check_two_phase_series_scan(use_index, use_range_index).await; + } + } +} + +async fn check_two_phase_series_scan(use_index: bool, use_range_index: bool) { let mut env = TestEnv::with_prefix("test_two_phase_series_scan").await; let engine = env .create_engine(MitoConfig { experimental_series_scan_v2: true, + experimental_enable_series_index: true, ..Default::default() }) .await; @@ -1313,12 +1322,106 @@ async fn test_two_phase_series_scan() { .await .unwrap(); test_util::flush_region(&engine, region_id, None).await; + if use_index || use_range_index { + use datatypes::arrow::array::{BinaryArray, UInt8Array, UInt64Array}; + + use crate::series_index::{ + SeriesIndexEntry, SeriesIndexFileHandle, SeriesIndexVersion, SeriesIndexWriter, + SeriesIndexWriterOptions, series_index_channel, series_index_path, + }; + + let region = engine.find_region(region_id).unwrap(); + let store = region.series_index_store.clone().unwrap(); + let sequence = region.flushed_sequence(); + let entry = SeriesIndexEntry { + index_uuid: FileId::random(), + bucket_start: Timestamp::new_millisecond(0), + bucket_end: Timestamp::new_millisecond(2001), + source_file_ids: Vec::new(), + min_file_sequence: sequence, + max_file_sequence: sequence, + compaction_window_secs: 1, + window_sequences: Default::default(), + }; + let keys = [ + (10, 0, "a", "x"), + (10, u64::MAX, "b", "y"), + (20, 0, "c", "z"), + ] + .map(|(table, tsid, a, b)| new_sparse_primary_key(&[a, b], &metadata, table, tsid)); + let batch = DfRecordBatch::try_from_iter(vec![ + ( + "ts", + Arc::new(TimestampMillisecondArray::from(vec![1000; 3])) as ArrayRef, + ), + ( + "__primary_key", + Arc::new(BinaryArray::from_iter_values(&keys)), + ), + ("__sequence", Arc::new(UInt64Array::from_value(sequence, 3))), + ("__op_type", Arc::new(UInt8Array::from_value(0, 3))), + ]) + .unwrap(); + let mut index_version = SeriesIndexVersion::default(); + if use_index { + let path = series_index_path(region_id, entry.index_uuid); + let mut writer = SeriesIndexWriter::try_new( + metadata.clone(), + store.clone(), + &path, + SeriesIndexWriterOptions::default(), + None, + ) + .await + .unwrap(); + writer.write(&batch).await.unwrap(); + writer.finish().await.unwrap(); + let (purger, _receiver) = series_index_channel(store.clone()); + let handle = SeriesIndexFileHandle::new(region_id, entry.clone(), purger); + index_version + .series_indexes + .insert(entry.index_uuid, handle); + } + if use_range_index { + use crate::sst::range_index::{ + SstRangeIndexWriter, SstRangeIndexWriterOptions, range_index_path, + }; + let version = region.version(); + let file = version + .ssts + .levels() + .iter() + .flat_map(|level| level.files.values()) + .next() + .unwrap(); + assert_eq!(1, file.meta_ref().num_row_groups); + let file_id = file.file_id().file_id(); + let mut writer = SstRangeIndexWriter::try_new( + metadata.clone(), + store, + &range_index_path(region_id, file_id), + SstRangeIndexWriterOptions::default(), + ) + .await + .unwrap(); + writer.write(0, &batch).await.unwrap(); + writer.finish().await.unwrap(); + index_version.range_indexes.insert(file_id); + } + region + .series_index_version_control + .publish(Arc::new(index_version)); + } + // A newer SST remains uncovered; the remaining writes stay in memory. + engine + .handle_request(region_id, put(rows(&[(10, 0, "a", "x", 11, 1000)]))) + .await + .unwrap(); + test_util::flush_region(&engine, region_id, None).await; engine .handle_request( region_id, put(rows(&[ - // Replaces the flushed row for this series and timestamp. - (10, 0, "a", "x", 11, 1000), (10, 0, "a", "x", 12, 2000), (10, u64::MAX, "b", "y", 21, 2000), (20, u64::MAX, "d", "w", 40, 1000), @@ -1370,6 +1473,12 @@ async fn test_two_phase_series_scan() { .await .unwrap(); + let index_files = metrics_set + .clone_inner() + .sum_by_name("candidate_index_files") + .map_or(0, |value| value.as_usize()); + assert_eq!(usize::from(use_index), index_files); + let mut series_to_partition = BTreeMap::new(); let mut actual_rows = Vec::new(); for (partition, batches) in partition_batches.into_iter().enumerate() { @@ -1457,6 +1566,65 @@ async fn test_two_phase_series_scan() { } replay_rows.sort(); assert_eq!(actual_rows, replay_rows); + + // Exercise precise field/time filtering and candidate tag filtering on both paths. + let filtered = engine + .scanner( + region_id, + ScanRequest { + projection: Some(vec![2, 4, 5]), + filters: vec![ + col("tag_0").eq(lit("a")), + col("field_0").eq(lit(11_u64)), + col("ts").eq(lit(ScalarValue::TimestampMillisecond(Some(1000), None))), + ], + distribution: Some(TimeSeriesDistribution::PerSeries), + ..Default::default() + }, + ) + .await + .unwrap(); + let batches = filtered + .scan() + .await + .unwrap() + .try_collect::>() + .await + .unwrap(); + assert_eq!( + 1, + batches.iter().map(|batch| batch.num_rows()).sum::() + ); + + if use_range_index { + // A fresh query must read the cataloged index: missing files must not + // silently switch back to the primary-key path. + let region = engine.find_region(region_id).unwrap(); + let version = region.series_index_version_control.current(); + let file_id = *version.range_indexes.iter().next().unwrap(); + region + .series_index_store + .as_ref() + .unwrap() + .delete(&crate::sst::range_index::range_index_path( + region_id, file_id, + )) + .await + .unwrap(); + let scanner = engine + .scanner( + region_id, + ScanRequest { + distribution: Some(TimeSeriesDistribution::PerSeries), + ..Default::default() + }, + ) + .await + .unwrap(); + if let Ok(stream) = scanner.scan().await { + assert!(stream.try_collect::>().await.is_err()); + } + } } /// Scans all partitions in round-robin fashion and returns rows sorted by (tag, ts). diff --git a/src/mito2/src/read/pruner.rs b/src/mito2/src/read/pruner.rs index 51630bbc62a..19c487ed0b8 100644 --- a/src/mito2/src/read/pruner.rs +++ b/src/mito2/src/read/pruner.rs @@ -87,6 +87,12 @@ impl PartitionPruner { } } + /// Excludes files replaced by another candidate source from prefetching. + pub(crate) fn excluding_files(mut self, excluded: &HashSet) -> Self { + self.file_indices.retain(|index| !excluded.contains(index)); + self + } + /// Gets or creates the FileRangeBuilder for a file. /// /// This method also triggers pre-fetching of upcoming files in the background @@ -737,7 +743,7 @@ impl Pruner { #[cfg(test)] impl Pruner { /// Returns the remaining range count for a file (test-only). - fn test_remaining_ranges(&self, file_index: usize) -> usize { + pub(crate) fn test_remaining_ranges(&self, file_index: usize) -> usize { self.inner.file_entries[file_index] .lock() .unwrap() diff --git a/src/mito2/src/read/scan_region.rs b/src/mito2/src/read/scan_region.rs index ef4c44a96cc..13864b05c24 100644 --- a/src/mito2/src/read/scan_region.rs +++ b/src/mito2/src/read/scan_region.rs @@ -76,6 +76,7 @@ use crate::read::unordered_scan::UnorderedScan; use crate::read::{BoxedRecordBatchStream, RecordBatch}; use crate::region::options::MergeMode; use crate::region::version::VersionRef; +use crate::series_index::SeriesIndexReadContext; use crate::sst::file::FileHandle; use crate::sst::index::bloom_filter::applier::{ BloomFilterIndexApplierBuilder, BloomFilterIndexApplierRef, @@ -235,6 +236,8 @@ impl Scanner { pub(crate) struct ScanRegion { /// Version of the region at scan. version: VersionRef, + /// Pinned index snapshot and its storage for candidate discovery. + series_index: Option, /// Access layer of the region. access_layer: AccessLayerRef, /// Scan request. @@ -276,6 +279,7 @@ impl ScanRegion { ) -> ScanRegion { ScanRegion { version, + series_index: None, access_layer, request, cache_strategy, @@ -294,6 +298,16 @@ impl ScanRegion { } } + /// Pins the series-index snapshot used by candidate discovery. + #[must_use] + pub(crate) fn with_series_index( + mut self, + series_index: Option, + ) -> Self { + self.series_index = series_index; + self + } + /// Sets counters that should receive query-load metrics. #[must_use] pub(crate) fn with_query_stat_counters(mut self, counters: RegionQueryStatCounters) -> Self { @@ -559,6 +573,7 @@ impl ScanRegion { }); let input = ScanInput::builder(self.access_layer, mapper) + .with_series_index(self.series_index) .with_time_range(Some(time_range)) .with_predicate(predicate) .with_memtables(mem_range_builders) @@ -927,6 +942,8 @@ fn time_range_covers_file(time_range: Option<&TimestampRange>, file: &FileHandle /// Common input for different scanners. pub struct ScanInput { + /// Pinned series-index snapshot and its storage, when configured. + pub(crate) series_index: Option, /// Region SST access layer. access_layer: AccessLayerRef, /// Maps projected Batches to RecordBatches. @@ -1025,6 +1042,7 @@ impl ScanInput { ) -> ScanInputBuilder { ScanInputBuilder { input: ScanInput { + series_index: None, access_layer, read_cols: mapper.read_columns().clone(), mapper: Arc::new(mapper), @@ -1089,6 +1107,16 @@ impl ScanInput { } impl ScanInputBuilder { + /// Sets the pinned series-index context for candidate discovery. + #[must_use] + pub(crate) fn with_series_index( + mut self, + series_index: Option, + ) -> Self { + self.input.series_index = series_index; + self + } + /// Sets time range filter for time index. #[must_use] pub(crate) fn with_time_range(mut self, time_range: Option) -> Self { @@ -1549,6 +1577,7 @@ impl ScanInput { let reader = self .access_layer .read_sst(file.clone()) + .series_index(self.series_index.clone()) .predicate(predicate) .projection(Some(self.read_cols.clone())) .json2_rewrite_targets(self.json2_rewrite_targets.clone()) diff --git a/src/mito2/src/read/series_candidate.rs b/src/mito2/src/read/series_candidate.rs index 1a5c7ca60a4..f05f54e6662 100644 --- a/src/mito2/src/read/series_candidate.rs +++ b/src/mito2/src/read/series_candidate.rs @@ -14,6 +14,8 @@ //! Candidate metric-series discovery for the two-stage series scan. +use std::cmp::Reverse; +use std::collections::HashSet; use std::sync::Arc; use std::time::Instant; @@ -21,7 +23,7 @@ use async_stream::try_stream; use datafusion::execution::memory_pool::{MemoryConsumer, MemoryPool}; use datafusion::physical_expr::{LexOrdering, PhysicalSortExpr}; use datafusion::physical_plan::expressions::Column; -use datafusion::physical_plan::metrics::{BaselineMetrics, ExecutionPlanMetricsSet}; +use datafusion::physical_plan::metrics::{BaselineMetrics, ExecutionPlanMetricsSet, MetricBuilder}; use datafusion::physical_plan::sorts::streaming_merge::StreamingMergeBuilder; use datafusion::physical_plan::stream::RecordBatchStreamAdapter; use datafusion_common::DataFusionError; @@ -50,7 +52,10 @@ use crate::read::range_cache::{ }; use crate::read::scan_region::StreamContext; use crate::read::scan_util::{PartitionMetrics, new_filter_metrics, scan_flat_mem_ranges}; -use crate::series_index::{METRIC_SERIES_ID_BATCH_SIZE, MetricSeriesId, MetricSeriesIdStream}; +use crate::series_index::{ + METRIC_SERIES_ID_BATCH_SIZE, MetricSeriesId, MetricSeriesIdStream, SeriesIndexFileHandle, + SeriesIndexReadContext, SeriesIndexSearcher, +}; use crate::sst::parquet::DEFAULT_READ_BATCH_SIZE; use crate::sst::parquet::format::PrimaryKeyArray; use crate::sst::parquet::prefilter::{ @@ -64,6 +69,8 @@ pub(crate) struct SeriesCandidateScanner { stream_ctx: Arc, partitions: Vec>, partition_pruner: Arc, + candidate_pruner: Arc, + coverage: Arc, range_semaphore: Arc, memory_pool: Arc, metrics_set: ExecutionPlanMetricsSet, @@ -108,11 +115,21 @@ impl SeriesCandidateScanner { ); let all_ranges = partitions.iter().flatten().copied().collect::>(); pruner.add_partition_ranges(&all_ranges); - let partition_pruner = Arc::new(PartitionPruner::new(pruner, &all_ranges)); + let coverage = Arc::new(SeriesIndexCoverage::new(&stream_ctx, &all_ranges)); + let partition_pruner = Arc::new(PartitionPruner::new(pruner.clone(), &all_ranges)); + let candidate_pruner = if coverage.covered_files.is_empty() { + partition_pruner.clone() + } else { + Arc::new( + PartitionPruner::new(pruner, &all_ranges).excluding_files(&coverage.covered_files), + ) + }; Ok(Self { stream_ctx, partitions, partition_pruner, + candidate_pruner, + coverage, range_semaphore, memory_pool, metrics_set, @@ -130,7 +147,8 @@ impl SeriesCandidateScanner { .collect::>(); let range_builder = SeriesCandidateRangeBuilder { stream_ctx: self.stream_ctx.clone(), - partition_pruner: self.partition_pruner.clone(), + partition_pruner: self.candidate_pruner.clone(), + coverage: self.coverage.clone(), range_semaphore: self.range_semaphore.clone(), memory_pool: self.memory_pool.clone(), metrics_set: self.metrics_set.clone(), @@ -162,6 +180,23 @@ impl SeriesCandidateScanner { range_streams.push(task.await.context(JoinSnafu)??); } + if let Some(context) = &self.stream_ctx.input.series_index { + MetricBuilder::new(&self.metrics_set) + .counter("candidate_index_files", self.partitions.len()) + .add(self.coverage.indexes.len()); + MetricBuilder::new(&self.metrics_set) + .counter("candidate_index_covered_ssts", self.partitions.len()) + .add(self.coverage.covered_files.len()); + for index in &self.coverage.indexes { + range_streams.push(index_primary_key_stream( + self.stream_ctx.clone(), + context.clone(), + index.clone(), + self.range_semaphore.clone(), + )); + } + } + // Keep scanner-level merge metrics in the same synthetic partition as // SeriesDistributor. Output partitions occupy 0..self.partitions.len(). let merged = merge_primary_key_streams( @@ -180,9 +215,137 @@ impl SeriesCandidateScanner { } } +/// A scanner-wide replacement plan, independent of partition-range boundaries. +#[derive(Default)] +struct SeriesIndexCoverage { + indexes: Vec, + /// Indices into `ScanInput.files`, shared by every occurrence of an SST. + covered_files: HashSet, +} + +impl SeriesIndexCoverage { + fn new(stream_ctx: &StreamContext, ranges: &[PartitionRange]) -> Self { + let Some(context) = &stream_ctx.input.series_index else { + return Self::default(); + }; + let mut uncovered: HashSet<_> = ranges + .iter() + .flat_map(|range| { + stream_ctx.ranges[range.identifier] + .row_group_indices + .iter() + .filter(|index| stream_ctx.is_file_range_index(**index)) + .map(|index| index.index - stream_ctx.input.num_memtables()) + }) + .collect(); + let region_id = stream_ctx.input.region_metadata().region_id; + let mut candidates: Vec<_> = context + .version + .series_indexes + .values() + .map(|index| { + let files: HashSet<_> = uncovered + .iter() + .copied() + .filter(|file_index| { + index + .entry() + .covers_file(stream_ctx.input.files[*file_index].meta_ref(), region_id) + }) + .collect(); + (index, files) + }) + .collect(); + let mut coverage = Self::default(); + while let Some((position, count)) = candidates + .iter() + .enumerate() + .map(|(position, (index, files))| { + ( + position, + files.intersection(&uncovered).count(), + index.entry(), + ) + }) + .max_by_key(|(_, count, entry)| { + ( + *count, + entry.max_file_sequence, + Reverse(entry.index_uuid.as_bytes()), + ) + }) + .map(|(position, count, _)| (position, count)) + { + if count == 0 { + break; + } + let (index, files) = candidates.swap_remove(position); + for file in files { + if uncovered.remove(&file) { + coverage.covered_files.insert(file); + } + } + coverage.indexes.push(index.clone()); + } + coverage + } + + fn covers_source(&self, stream_ctx: &StreamContext, index: RowGroupIndex) -> bool { + stream_ctx.is_file_range_index(index) + && self + .covered_files + .contains(&(index.index - stream_ctx.input.num_memtables())) + } +} + +/// Reads an index once for the entire scan, rather than once per partition range. +fn index_primary_key_stream( + stream_ctx: Arc, + context: SeriesIndexReadContext, + index: SeriesIndexFileHandle, + semaphore: Arc, +) -> BoxedRecordBatchStream { + Box::pin(try_stream! { + let metadata = stream_ctx.input.region_metadata(); + let codec = SparsePrimaryKeyCodec::new(metadata); + let mut series = { + let _permit = semaphore.acquire().await.map_err(|error| UnexpectedSnafu { + reason: format!("failed to acquire candidate index permit: {error}"), + }.build())?; + SeriesIndexSearcher::try_new( + metadata.clone(), + context.store.clone(), + index, + stream_ctx.input.predicate_group().predicate(), + stream_ctx.input.time_range, + ).await?.search()? + }; + loop { + let batch = { + let _permit = semaphore.acquire().await.map_err(|error| UnexpectedSnafu { + reason: format!("failed to acquire candidate index permit: {error}"), + }.build())?; + series.try_next().await? + }; + let Some(batch) = batch else { break }; + let mut builder = BinaryBuilder::new(); + let mut key = Vec::new(); + for series in batch { + key.clear(); + codec.encode_internal(series.table_id, series.tsid, &mut key) + .context(crate::error::EncodeSnafu)?; + builder.append_value(&key); + } + yield RecordBatch::try_new(primary_key_schema(), vec![Arc::new(builder.finish())]) + .context(NewRecordBatchSnafu)?; + } + }) +} + #[derive(Clone)] struct SeriesCandidateRangeBuilder { stream_ctx: Arc, + coverage: Arc, partition_pruner: Arc, range_semaphore: Arc, memory_pool: Arc, @@ -196,7 +359,18 @@ impl SeriesCandidateRangeBuilder { part_range: PartitionRange, merge_partition: usize, ) -> Result { - let cache_key = build_candidate_range_cache_key(&self.stream_ctx, &part_range); + let range_meta = &self.stream_ctx.ranges[part_range.identifier]; + // A cache entry describes the complete original range. Never cache a + // partial range whose missing candidates are supplied by a global index. + let replaced_sources = range_meta + .row_group_indices + .iter() + .any(|index| self.coverage.covers_source(&self.stream_ctx, *index)); + let cache_key = if replaced_sources { + None + } else { + build_candidate_range_cache_key(&self.stream_ctx, &part_range) + }; if let Some(key) = cache_key.as_ref() { if let Some(value) = self.stream_ctx.input.cache_strategy.get_range_result(key) { self.part_metrics.inc_range_cache_hit(); @@ -205,7 +379,6 @@ impl SeriesCandidateRangeBuilder { self.part_metrics.inc_range_cache_miss(); } - let range_meta = &self.stream_ctx.ranges[part_range.identifier]; let mut sources = Vec::with_capacity(range_meta.row_group_indices.len()); for index in &range_meta.row_group_indices { let source = self.build_source(*index, range_meta.time_range).await?; @@ -260,6 +433,12 @@ impl SeriesCandidateRangeBuilder { } if self.stream_ctx.is_file_range_index(index) { + if self.coverage.covers_source(&self.stream_ctx, index) { + // Leave the range reference for the data phase so its first + // read can cache the builder. Retaining builders only prevents + // eviction; it does not allow caching at zero references. + return Ok(None); + } let file = self.stream_ctx.input.file_from_index(index); let predicate = self.stream_ctx.input.predicate_for_file(file); if self @@ -561,23 +740,260 @@ fn decode_metric_series( #[cfg(test)] mod tests { + use std::num::NonZeroU64; use std::time::Instant; + use common_time::Timestamp; use datafusion::execution::memory_pool::UnboundedMemoryPool; use datafusion_expr::{col, lit}; - use datatypes::arrow::array::{ArrayRef, DictionaryArray, UInt32Array}; + use datatypes::arrow::array::{ + ArrayRef, DictionaryArray, TimestampMillisecondArray, UInt8Array, UInt32Array, UInt64Array, + }; use datatypes::arrow::datatypes::UInt32Type; use futures::TryStreamExt; use store_api::codec::PrimaryKeyEncoding; + use store_api::storage::FileId; use table::predicate::Predicate; use super::*; + use crate::cache::{CacheManager, CacheStrategy}; use crate::read::flat_projection::FlatProjectionMapper; + use crate::read::pruner::PrunerOptions; use crate::read::scan_region::ScanInput; use crate::read::scan_util::PartitionMetrics; + use crate::series_index::{ + SeriesIndexEntry, SeriesIndexVersion, SeriesIndexWriter, SeriesIndexWriterOptions, + series_index_channel, series_index_path, + }; + use crate::sst::file::{FileHandle, FileMeta}; + use crate::test_util::new_noop_file_purger; use crate::test_util::scheduler_util::SchedulerEnv; use crate::test_util::sst_util::sst_region_metadata_with_encoding; + /// Uses absent SST objects so any covered-SST read fails the test. + async fn indexed_scanner() -> (SchedulerEnv, SeriesCandidateScanner, Arc) { + let env = SchedulerEnv::new().await; + let metadata = Arc::new(sst_region_metadata_with_encoding( + PrimaryKeyEncoding::Sparse, + )); + let store = env.access_layer.object_store().clone(); + let entry = SeriesIndexEntry { + index_uuid: FileId::random(), + bucket_start: Timestamp::new_millisecond(0), + bucket_end: Timestamp::new_millisecond(20), + source_file_ids: Vec::new(), + min_file_sequence: 1, + max_file_sequence: 2, + compaction_window_secs: 1, + window_sequences: Default::default(), + }; + let path = series_index_path(metadata.region_id, entry.index_uuid); + let codec = SparsePrimaryKeyCodec::new(&metadata); + let keys: Vec<_> = (0..1001) + .map(|tsid| { + let mut key = Vec::new(); + codec.encode_internal(1, tsid, &mut key).unwrap(); + key + }) + .collect(); + let batch = RecordBatch::try_from_iter(vec![ + ( + "ts", + Arc::new(TimestampMillisecondArray::from(vec![10; keys.len()])) as ArrayRef, + ), + ( + "__primary_key", + Arc::new(BinaryArray::from_iter_values(&keys)), + ), + ( + "__sequence", + Arc::new(UInt64Array::from_value(1, keys.len())), + ), + ("__op_type", Arc::new(UInt8Array::from_value(0, keys.len()))), + ]) + .unwrap(); + let mut writer = SeriesIndexWriter::try_new( + metadata.clone(), + store.clone(), + &path, + SeriesIndexWriterOptions { + row_group_size: 500, + }, + None, + ) + .await + .unwrap(); + writer.write(&batch).await.unwrap(); + writer.finish().await.unwrap(); + let (purger, _receiver) = series_index_channel(store.clone()); + let handle = SeriesIndexFileHandle::new(metadata.region_id, entry.clone(), purger); + let context = SeriesIndexReadContext { + store, + version: Arc::new(SeriesIndexVersion { + series_indexes: [(entry.index_uuid, handle)].into(), + ..Default::default() + }), + }; + let files = (0..2) + .map(|i| { + FileHandle::new( + FileMeta { + region_id: metadata.region_id, + file_id: FileId::random(), + time_range: ( + Timestamp::new_millisecond(i * 10), + Timestamp::new_millisecond(i * 10 + 9), + ), + sequence: NonZeroU64::new(i as u64 + 1), + num_row_groups: 1, + ..Default::default() + }, + new_noop_file_purger(), + ) + }) + .collect(); + let mapper = + FlatProjectionMapper::new(&metadata, 0..metadata.column_metadatas.len()).unwrap(); + let input = ScanInput::builder(env.access_layer.clone(), mapper) + .with_predicate( + crate::read::scan_region::PredicateGroup::new( + &metadata, + &[col("__table_id").eq(lit(1_u32))], + ) + .unwrap(), + ) + .with_files(files) + .with_series_index(Some(context)) + .with_cache(CacheStrategy::EnableAll(Arc::new( + CacheManager::builder() + .range_result_cache_size(1024 * 1024) + .build(), + ))) + .build(); + let stream_ctx = Arc::new(StreamContext::seq_scan_ctx(input)); + let ranges = stream_ctx.partition_ranges(); + assert_eq!(2, ranges.len()); + let metrics_set = ExecutionPlanMetricsSet::new(); + let part_metrics = PartitionMetrics::new( + metadata.region_id, + 2, + "candidate-test", + Instant::now(), + false, + &metrics_set, + ); + let pruner = Arc::new(Pruner::new_with_options( + stream_ctx.clone(), + 1, + PrunerOptions { + retain_builders: true, + enable_predicate_prefilter: false, + }, + )); + let scanner = SeriesCandidateScanner::try_new( + stream_ctx, + ranges.into_iter().map(|range| vec![range]).collect(), + pruner.clone(), + Arc::new(Semaphore::new(1)), + Arc::new(UnboundedMemoryPool::default()), + metrics_set, + part_metrics, + ) + .unwrap(); + (env, scanner, pruner) + } + + #[tokio::test] + async fn index_is_shared_across_ranges_without_caching_partial_candidates() { + let (_env, scanner, pruner) = indexed_scanner().await; + let keys: Vec<_> = scanner + .partitions + .iter() + .flatten() + .map(|range| build_candidate_range_cache_key(&scanner.stream_ctx, range).unwrap()) + .collect(); + let groups = scanner + .build_stream() + .await + .unwrap() + .try_collect::>() + .await + .unwrap(); + assert_eq!( + (0..1001) + .map(|tsid| MetricSeriesId { table_id: 1, tsid }) + .collect::>(), + groups.into_iter().flatten().collect::>() + ); + // Covered SSTs have no builder yet. Keep their references so the data + // phase can cache each builder on its first read and reuse it later. + for file_index in 0..2 { + assert_eq!(1, pruner.test_remaining_ranges(file_index)); + } + assert_eq!( + 1, + scanner + .metrics_set + .clone_inner() + .sum_by_name("candidate_index_files") + .unwrap() + .as_usize() + ); + assert_eq!( + 2, + scanner + .metrics_set + .clone_inner() + .sum_by_name("candidate_index_covered_ssts") + .unwrap() + .as_usize() + ); + // Index-backed ranges must not publish their empty residual streams as + // complete candidate results for a future scan without the index. + assert!(!format!("{:?}", scanner.part_metrics).contains("range_cache_miss")); + for key in keys { + assert!( + scanner + .stream_ctx + .input + .cache_strategy + .get_range_result(&key) + .is_none() + ); + } + } + + #[tokio::test] + async fn index_read_failure_after_output_is_propagated() { + let (_env, scanner, _pruner) = indexed_scanner().await; + let context = scanner.stream_ctx.input.series_index.clone().unwrap(); + let index = scanner.coverage.indexes[0].clone(); + let path = series_index_path( + scanner.stream_ctx.input.region_metadata().region_id, + index.entry().index_uuid, + ); + let mut stream = index_primary_key_stream( + scanner.stream_ctx.clone(), + context.clone(), + index, + scanner.range_semaphore.clone(), + ); + assert_eq!(500, stream.try_next().await.unwrap().unwrap().num_rows()); + context.store.delete(&path).await.unwrap(); + assert!(stream.try_collect::>().await.is_err()); + // Opening the missing selected index must also fail, rather than emit + // an incomplete candidate set or fall back to the absent SSTs. + assert!( + scanner + .build_stream() + .await + .unwrap() + .try_collect::>() + .await + .is_err() + ); + } + #[tokio::test] async fn candidate_scanner_rejects_predicate_prefilter_pruner() { let env = SchedulerEnv::new().await; diff --git a/src/mito2/src/read/series_reader.rs b/src/mito2/src/read/series_reader.rs index 1164905afa7..9e86ff1e5e9 100644 --- a/src/mito2/src/read/series_reader.rs +++ b/src/mito2/src/read/series_reader.rs @@ -159,15 +159,20 @@ impl SeriesBatchCollector { struct MetricSeriesFilter { range: SeriesRange, series: Arc>, + sorted_series: Arc>, enable_range_cache: bool, } impl MetricSeriesFilter { fn new(assigned: &AssignedSeriesBatch) -> Self { - let series = assigned.series().iter().copied().collect(); + let mut sorted_series = assigned.series().to_vec(); + sorted_series.sort_unstable(); + sorted_series.dedup(); + let series = sorted_series.iter().copied().collect(); Self { range: assigned.range(), series: Arc::new(series), + sorted_series: Arc::new(sorted_series), enable_range_cache: assigned.enable_range_cache(), } } @@ -256,7 +261,6 @@ fn filter_flat_stream_by_series( pub(crate) struct SeriesReader { stream_ctx: Arc, partition_ranges: Vec, - range: SeriesRange, filter: MetricSeriesFilter, codec: SparsePrimaryKeyCodec, partition_pruner: Arc, @@ -292,13 +296,11 @@ impl SeriesReader { } ); - let range = assigned_series.range(); let filter = MetricSeriesFilter::new(&assigned_series); let codec = SparsePrimaryKeyCodec::new(stream_ctx.input.region_metadata()); Ok(Self { stream_ctx, partition_ranges, - range, filter, codec, partition_pruner, @@ -320,7 +322,6 @@ impl SeriesReader { let partition_pruner = self.partition_pruner.clone(); let range_semaphore = self.range_semaphore.clone(); let part_metrics = self.part_metrics.clone(); - let range = self.range; tasks.push(common_runtime::spawn_query(async move { let _permit = range_semaphore.acquire().await.map_err(|error| { UnexpectedSnafu { @@ -331,7 +332,6 @@ impl SeriesReader { build_series_partition_range( stream_ctx, part_range, - range, filter, codec, partition_pruner, @@ -367,12 +367,12 @@ impl SeriesReader { async fn build_series_partition_range( stream_ctx: Arc, part_range: PartitionRange, - range: SeriesRange, filter: MetricSeriesFilter, codec: SparsePrimaryKeyCodec, partition_pruner: Arc, part_metrics: PartitionMetrics, ) -> Result<(BoxedRecordBatchStream, usize)> { + let range = filter.range; let cache_key = filter .enable_range_cache .then(|| build_series_range_cache_key(&stream_ctx, &part_range, range)) @@ -512,18 +512,22 @@ fn scan_series_file_ranges( for range in ranges { let build_start = Instant::now(); - let Some(mut reader) = range - .reader_by_primary_key( - primary_key_filter.as_mut(), - fetch_metrics.as_deref(), - ) - .await? - else { - continue; + let searcher = range.range_index_searcher().await?; + let reader = if let Some(searcher) = searcher { + let row_group_id = u32::try_from(range.row_group_index()).map_err(|_| UnexpectedSnafu { + reason: format!("row group index exceeds u32: {}", range.row_group_index()), + }.build())?; + let selected = searcher.search(row_group_id, &filter.sorted_series).await?; + range.reader_by_row_ranges(selected, fetch_metrics.as_deref()).await? + } else { + range.reader_by_primary_key(primary_key_filter.as_mut(), fetch_metrics.as_deref()).await? }; let build_cost = build_start.elapsed(); reader_metrics.build_cost += build_cost; part_metrics.inc_build_reader_cost(build_cost); + let Some(mut reader) = reader else { + continue; + }; let scan_start = Instant::now(); let file_sequence_trusted = range diff --git a/src/mito2/src/region.rs b/src/mito2/src/region.rs index bdf18550cd6..dd392f5b0ab 100644 --- a/src/mito2/src/region.rs +++ b/src/mito2/src/region.rs @@ -29,6 +29,7 @@ use common_base::hash::partition_expr_version; use common_recordbatch::adapter::RegionQueryStatCounters; use common_telemetry::{error, info, warn}; use crossbeam_utils::atomic::AtomicCell; +use object_store::ObjectStore; use partition::expr::PartitionExpr; use snafu::{OptionExt, ResultExt, ensure}; use store_api::ManifestVersion; @@ -153,6 +154,8 @@ pub struct MitoRegion { pub(crate) version_control: VersionControlRef, /// Snapshot controller for range and series indexes. pub(crate) series_index_version_control: SeriesIndexVersionControl, + /// Store containing the region's series indexes. + pub(crate) series_index_store: Option, /// SSTs accessor for this region. pub(crate) access_layer: AccessLayerRef, /// Context to maintain manifest for this region. @@ -2049,6 +2052,7 @@ mod tests { region_id: metadata.region_id, version_control, series_index_version_control: Default::default(), + series_index_store: None, access_layer: env.access_layer.clone(), manifest_ctx, file_purger: crate::test_util::new_noop_file_purger(), @@ -2553,6 +2557,7 @@ mod tests { region_id: metadata.region_id, version_control, series_index_version_control: Default::default(), + series_index_store: None, access_layer, manifest_ctx: manifest_ctx.clone(), file_purger: crate::test_util::new_noop_file_purger(), diff --git a/src/mito2/src/region/opener.rs b/src/mito2/src/region/opener.rs index ff5d3999f17..033f9916905 100644 --- a/src/mito2/src/region/opener.rs +++ b/src/mito2/src/region/opener.rs @@ -440,6 +440,7 @@ impl RegionOpener { region_id, version_control, series_index_version_control: Default::default(), + series_index_store: self.series_index_store.clone(), access_layer: access_layer.clone(), // Region is writable after it is created. manifest_ctx: Arc::new(ManifestContext::new( @@ -686,6 +687,7 @@ impl RegionOpener { region_id: self.region_id, version_control: version_control.clone(), series_index_version_control, + series_index_store: self.series_index_store.clone(), access_layer: access_layer.clone(), // Region is always opened in read only mode. manifest_ctx: Arc::new(ManifestContext::new( diff --git a/src/mito2/src/series_index.rs b/src/mito2/src/series_index.rs index 6982e4398f4..4a4133e1a9b 100644 --- a/src/mito2/src/series_index.rs +++ b/src/mito2/src/series_index.rs @@ -28,22 +28,38 @@ mod tests; mod version; mod writer; +use std::sync::Arc; + use futures::stream::BoxStream; +use object_store::ObjectStore; use store_api::metric_engine_consts::{ DATA_SCHEMA_TABLE_ID_COLUMN_NAME as TABLE_ID_COLUMN, DATA_SCHEMA_TSID_COLUMN_NAME as TSID_COLUMN, }; use crate::error::Result; -pub(crate) use crate::series_index::catalog::{delete_catalogs, load_version_control}; +#[cfg(test)] +pub(crate) use crate::series_index::catalog::SeriesIndexEntry; +pub(crate) use crate::series_index::catalog::{ + delete_catalogs, load_version_control, series_index_path, +}; pub(crate) use crate::series_index::purger::{IndexFilePurger, series_index_channel}; pub use crate::series_index::searcher::SeriesIndexSearcher; pub(crate) use crate::series_index::task::{SeriesIndexTaskState, spawn_series_index_tasks}; -pub(crate) use crate::series_index::version::{SeriesIndexVersion, SeriesIndexVersionControl}; +pub(crate) use crate::series_index::version::{ + SeriesIndexFileHandle, SeriesIndexVersion, SeriesIndexVersionControl, +}; pub use crate::series_index::writer::{ SeriesIndexWriter, SeriesIndexWriterMetrics, SeriesIndexWriterOptions, series_index_schema, }; +/// Index storage and pinned catalog snapshot for a query. +#[derive(Clone)] +pub(crate) struct SeriesIndexReadContext { + pub(crate) store: ObjectStore, + pub(crate) version: Arc, +} + pub(crate) const MIN_TS_COLUMN: &str = "__series_min_ts"; pub(crate) const MAX_TS_COLUMN: &str = "__series_max_ts"; pub(crate) const ROW_COUNT_COLUMN: &str = "__series_row_count"; diff --git a/src/mito2/src/series_index/catalog.rs b/src/mito2/src/series_index/catalog.rs index 56c61da54f9..d56e10698fd 100644 --- a/src/mito2/src/series_index/catalog.rs +++ b/src/mito2/src/series_index/catalog.rs @@ -54,6 +54,10 @@ pub(crate) struct WindowSequence { } /// Self-describing coverage stored in a series-index Parquet footer. +/// +/// A published entry must include every series from every SST in the region +/// contained by its bucket and inclusive file-sequence interval. Query planning +/// relies on this complete-coverage contract, not on `source_file_ids`. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub(crate) struct SeriesIndexEntry { pub(crate) index_uuid: FileId, @@ -76,6 +80,23 @@ pub(crate) struct SeriesIndexEntry { pub(crate) window_sequences: BTreeMap, } +impl SeriesIndexEntry { + /// Whether this index completely covers an SST in the query's sequence domain. + pub(crate) fn covers_file( + &self, + file: &crate::sst::file::FileMeta, + region_id: RegionId, + ) -> bool { + file.region_id == region_id + && file.time_range.0 <= file.time_range.1 + && self.bucket_start <= file.time_range.0 + && file.time_range.1 < self.bucket_end + && file.sequence.is_some_and(|sequence| { + self.min_file_sequence <= sequence.get() && sequence.get() <= self.max_file_sequence + }) + } +} + #[derive(Debug, Default, Serialize, Deserialize)] pub(crate) struct SeriesIndexCatalog { pub(crate) indexes: Vec, @@ -225,6 +246,46 @@ mod tests { } } + #[test] + fn coverage_uses_exclusive_time_end_and_inclusive_file_sequences() { + let region_id = RegionId::new(1, 1); + let entry = SeriesIndexEntry { + index_uuid: FileId::random(), + bucket_start: Timestamp::new_second(1), + bucket_end: Timestamp::new_second(2), + source_file_ids: Vec::new(), + min_file_sequence: 2, + max_file_sequence: 4, + compaction_window_secs: 1, + window_sequences: BTreeMap::new(), + }; + for (start, end, sequence, own_region, covered) in [ + (1000, 1999, 2, true, true), + (1000, 1999, 4, true, true), + (999, 1999, 3, true, false), + (1000, 2000, 3, true, false), + (1000, 1999, 1, true, false), + (1000, 1999, 5, true, false), + (1000, 1999, 0, true, false), + (1000, 1999, 3, false, false), + ] { + let file = crate::sst::file::FileMeta { + region_id: if own_region { + region_id + } else { + RegionId::new(2, 1) + }, + time_range: ( + Timestamp::new_millisecond(start), + Timestamp::new_millisecond(end), + ), + sequence: std::num::NonZeroU64::new(sequence), + ..Default::default() + }; + assert_eq!(covered, entry.covers_file(&file, region_id), "{file:?}"); + } + } + #[tokio::test] async fn test_load_catalog_defaults_on_missing_invalid_or_unreadable_catalog() { let store = ObjectStore::new(Memory::default()).unwrap(); diff --git a/src/mito2/src/series_index/searcher.rs b/src/mito2/src/series_index/searcher.rs index b5b45e2ed63..4613edae913 100644 --- a/src/mito2/src/series_index/searcher.rs +++ b/src/mito2/src/series_index/searcher.rs @@ -31,13 +31,16 @@ use table::predicate::Predicate; use crate::error::{InvalidRecordBatchSnafu, RecordBatchSnafu, Result, UnexpectedSnafu}; use crate::series_index::{ MAX_TS_COLUMN, METRIC_SERIES_ID_BATCH_SIZE, MIN_TS_COLUMN, MetricSeriesId, - MetricSeriesIdStream, ROW_COUNT_COLUMN, TABLE_ID_COLUMN, TSID_COLUMN, series_index_schema, + MetricSeriesIdStream, ROW_COUNT_COLUMN, SeriesIndexFileHandle, TABLE_ID_COLUMN, TSID_COLUMN, + series_index_path, series_index_schema, }; use crate::sst::parquet::index_reader::ParquetIndexReader; use crate::sst::parquet::prefilter::simple_tag_filters; /// Searches a series-index file for metric series matching query predicates. pub struct SeriesIndexSearcher { + /// Pins the index file until this searcher and all of its streams are dropped. + file_handle: SeriesIndexFileHandle, /// The file to search, or `None` when an empty time range rules out /// every series without opening the file. reader: Option, @@ -47,14 +50,14 @@ pub struct SeriesIndexSearcher { } impl SeriesIndexSearcher { - /// Creates a searcher for the series-index file of `metadata` at `path`. + /// Creates a searcher for the series-index file protected by `file_handle`. /// Predicates are built from the unit recorded in the file's schema, so a /// file written before a time index unit widening keeps being interpreted /// in its own unit. - pub async fn try_new( + pub(crate) async fn try_new( metadata: RegionMetadataRef, object_store: ObjectStore, - path: &str, + file_handle: SeriesIndexFileHandle, predicate: Option<&Predicate>, time_range: Option, ) -> Result { @@ -62,13 +65,16 @@ impl SeriesIndexSearcher { series_index_schema(&metadata)?; if time_range.as_ref().is_some_and(TimestampRange::is_empty) { return Ok(Self { + file_handle, reader: None, pruning_predicate: Predicate::new(Vec::new()), filters: Vec::new(), }); } - let reader = ParquetIndexReader::open(object_store, path).await?; + let file_id = file_handle.file_id(); + let path = series_index_path(file_id.region_id(), file_id.file_id()); + let reader = ParquetIndexReader::open(object_store, &path).await?; let unit = validate_index_schema(reader.schema())?; let mut filters = simple_tag_filters(&metadata, None, predicate); @@ -83,6 +89,7 @@ impl SeriesIndexSearcher { let (pruning_predicate, filters) = filters_for_schema(reader.schema(), &filters); Ok(Self { + file_handle, reader: Some(reader), pruning_predicate, filters, @@ -100,6 +107,7 @@ impl SeriesIndexSearcher { projection_columns.extend(self.filters.iter().map(SimpleFilterEvaluator::column_name)); let mut batches = reader.read(&self.pruning_predicate, &projection_columns)?; let filters = self.filters.clone(); + let file_handle = self.file_handle.clone(); Ok(Box::pin(try_stream! { let mut last_series = None; @@ -149,6 +157,7 @@ impl SeriesIndexSearcher { if !output.is_empty() { yield output; } + drop(file_handle); })) } } @@ -283,15 +292,41 @@ mod tests { use object_store::services::Memory; use store_api::codec::PrimaryKeyEncoding; use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder}; + use store_api::storage::FileId; + use tokio::sync::mpsc::UnboundedReceiver; use super::*; - use crate::series_index::{SeriesIndexWriter, SeriesIndexWriterOptions}; + use crate::series_index::purger::PurgeRequest; + use crate::series_index::{ + SeriesIndexEntry, SeriesIndexWriter, SeriesIndexWriterOptions, series_index_channel, + }; use crate::test_util::sst_util::{new_sparse_primary_key, sst_region_metadata_with_encoding}; fn object_store() -> ObjectStore { ObjectStore::new(Memory::default()).unwrap() } + fn index_handle( + metadata: &RegionMetadataRef, + store: &ObjectStore, + ) -> (SeriesIndexFileHandle, UnboundedReceiver) { + let (purger, receiver) = series_index_channel(store.clone()); + let entry = SeriesIndexEntry { + index_uuid: FileId::random(), + bucket_start: common_time::Timestamp::new_second(0), + bucket_end: common_time::Timestamp::new_second(60), + source_file_ids: Vec::new(), + min_file_sequence: 0, + max_file_sequence: 0, + compaction_window_secs: 60, + window_sequences: Default::default(), + }; + ( + SeriesIndexFileHandle::new(metadata.region_id, entry, purger), + receiver, + ) + } + fn flat_batch_with_time_unit( primary_keys: &[Vec], timestamps: &[i64], @@ -334,14 +369,12 @@ mod tests { async fn write_index( metadata: RegionMetadataRef, object_store: ObjectStore, - path: &str, rows: &[(u32, u64, &str, &str, i64)], row_group_size: usize, - ) { + ) -> (SeriesIndexFileHandle, UnboundedReceiver) { write_index_with_time_unit( metadata, object_store, - path, rows, row_group_size, ArrowTimeUnit::Millisecond, @@ -352,11 +385,13 @@ mod tests { async fn write_index_with_time_unit( metadata: RegionMetadataRef, object_store: ObjectStore, - path: &str, rows: &[(u32, u64, &str, &str, i64)], row_group_size: usize, unit: ArrowTimeUnit, - ) { + ) -> (SeriesIndexFileHandle, UnboundedReceiver) { + let (file_handle, receiver) = index_handle(&metadata, &object_store); + let file_id = file_handle.file_id(); + let path = series_index_path(file_id.region_id(), file_id.file_id()); let primary_keys = rows .iter() .map(|(table_id, tsid, tag_0, tag_1, _)| { @@ -367,7 +402,7 @@ mod tests { let mut writer = SeriesIndexWriter::try_new( metadata, object_store, - path, + &path, SeriesIndexWriterOptions { row_group_size }, None, ) @@ -378,6 +413,7 @@ mod tests { .await .unwrap(); writer.finish().await.unwrap(); + (file_handle, receiver) } async fn collect_ids(stream: MetricSeriesIdStream) -> Vec { @@ -396,11 +432,9 @@ mod tests { PrimaryKeyEncoding::Sparse, )); let object_store = object_store(); - let path = "search.parquet"; - write_index( + let (index, _receiver) = write_index( metadata.clone(), object_store.clone(), - path, &[ (1, 10, "a", "x", 10), (1, 20, "b", "x", 20), @@ -424,7 +458,7 @@ mod tests { let searcher = SeriesIndexSearcher::try_new( metadata.clone(), object_store.clone(), - path, + index.clone(), Some(&predicate), Some(time_range), ) @@ -449,7 +483,7 @@ mod tests { let searcher = SeriesIndexSearcher::try_new( metadata.clone(), object_store.clone(), - path, + index.clone(), None, Some(time_range), ) @@ -472,7 +506,7 @@ mod tests { ) .unwrap(); let searcher = - SeriesIndexSearcher::try_new(metadata, object_store, path, None, Some(time_range)) + SeriesIndexSearcher::try_new(metadata, object_store, index, None, Some(time_range)) .await .unwrap(); let ids = collect_ids(searcher.search().unwrap()).await; @@ -491,11 +525,9 @@ mod tests { PrimaryKeyEncoding::Sparse, )); let object_store = object_store(); - let path = "schema-evolution.parquet"; - write_index( + let (index, _receiver) = write_index( old_metadata.clone(), object_store.clone(), - path, &[ (1, 10, "a", "x", 10), (1, 20, "b", "x", 20), @@ -521,7 +553,7 @@ mod tests { let searcher = SeriesIndexSearcher::try_new( current_metadata, object_store, - path, + index, Some(&predicate), None, ) @@ -550,11 +582,9 @@ mod tests { PrimaryKeyEncoding::Sparse, )); let object_store = object_store(); - let path = "widen.parquet"; - write_index( + let (index, _receiver) = write_index( metadata.clone(), object_store.clone(), - path, &[ (1, 10, "a", "x", 10), (1, 20, "b", "x", 20), @@ -585,7 +615,7 @@ mod tests { let searcher = SeriesIndexSearcher::try_new( widened, object_store.clone(), - path, + index.clone(), None, Some(time_range), ) @@ -608,7 +638,7 @@ mod tests { ) .unwrap(); let searcher = - SeriesIndexSearcher::try_new(metadata, object_store, path, None, Some(time_range)) + SeriesIndexSearcher::try_new(metadata, object_store, index, None, Some(time_range)) .await .unwrap(); let ids = collect_ids(searcher.search().unwrap()).await; @@ -629,10 +659,9 @@ mod tests { PrimaryKeyEncoding::Sparse, )); let object_store = object_store(); - write_index( + let (index, _receiver) = write_index( metadata.clone(), object_store.clone(), - "old_ms.parquet", &[(1, 10, "a", "x", 10), (1, 20, "b", "x", 20)], 2, ) @@ -647,10 +676,9 @@ mod tests { } } let widened = Arc::new(widened); - write_index_with_time_unit( + let (new_index, _new_receiver) = write_index_with_time_unit( widened.clone(), object_store.clone(), - "new_us.parquet", &[(1, 30, "c", "x", 15_000), (1, 40, "d", "x", 15)], 2, ArrowTimeUnit::Microsecond, @@ -671,7 +699,7 @@ mod tests { let searcher = SeriesIndexSearcher::try_new( widened.clone(), object_store.clone(), - "old_ms.parquet", + index, None, Some(time_range), ) @@ -685,15 +713,10 @@ mod tests { tsid: 20 }] ); - let searcher = SeriesIndexSearcher::try_new( - widened, - object_store, - "new_us.parquet", - None, - Some(time_range), - ) - .await - .unwrap(); + let searcher = + SeriesIndexSearcher::try_new(widened, object_store, new_index, None, Some(time_range)) + .await + .unwrap(); let ids = collect_ids(searcher.search().unwrap()).await; assert_eq!( ids, @@ -761,19 +784,12 @@ mod tests { let rows = (0..501_u64) .map(|tsid| (1, tsid, "a", "x", tsid as i64)) .collect::>(); - write_index( - metadata.clone(), - object_store.clone(), - "batching.parquet", - &rows, - 100, - ) - .await; + let (index, _receiver) = + write_index(metadata.clone(), object_store.clone(), &rows, 100).await; - let searcher = - SeriesIndexSearcher::try_new(metadata, object_store, "batching.parquet", None, None) - .await - .unwrap(); + let searcher = SeriesIndexSearcher::try_new(metadata, object_store, index, None, None) + .await + .unwrap(); let batches = searcher .search() .unwrap() @@ -797,17 +813,54 @@ mod tests { ); } + #[tokio::test] + async fn search_streams_pin_deleted_index_until_released() { + let metadata = Arc::new(sst_region_metadata_with_encoding( + PrimaryKeyEncoding::Sparse, + )); + let store = object_store(); + let rows = (0..501_u64) + .map(|tsid| (1, tsid, "a", "x", tsid as i64)) + .collect::>(); + let (index, mut receiver) = write_index(metadata.clone(), store.clone(), &rows, 100).await; + let file_id = index.file_id(); + let searcher = SeriesIndexSearcher::try_new(metadata, store, index.clone(), None, None) + .await + .unwrap(); + index.mark_deleted(); + drop(index); + assert!(receiver.try_recv().is_err()); + + let mut stream = searcher.search().unwrap(); + let cancelled = searcher.search().unwrap(); + drop(searcher); + // Even a stream that has not been polled must pin its file. + assert!(receiver.try_recv().is_err()); + assert_eq!(stream.try_next().await.unwrap().unwrap().len(), 500); + assert!(receiver.try_recv().is_err()); + assert_eq!( + collect_ids(stream).await, + vec![MetricSeriesId { + table_id: 1, + tsid: 500 + }] + ); + assert!(receiver.try_recv().is_err()); + + // The unpolled stream is now the final owner; cancellation releases it. + drop(cancelled); + assert_eq!(receiver.try_recv().unwrap().file_id, file_id); + } + #[tokio::test] async fn search_prunes_row_groups_and_empty_ranges() { let metadata = Arc::new(sst_region_metadata_with_encoding( PrimaryKeyEncoding::Sparse, )); let object_store = object_store(); - let path = "pruning.parquet"; - write_index( + let (index, _receiver) = write_index( metadata.clone(), object_store.clone(), - path, &[ (1, 0, "a", "x", 0), (1, 1, "b", "x", 1), @@ -820,7 +873,9 @@ mod tests { ) .await; - let reader = ParquetIndexReader::open(object_store.clone(), path) + let file_id = index.file_id(); + let path = series_index_path(file_id.region_id(), file_id.file_id()); + let reader = ParquetIndexReader::open(object_store.clone(), &path) .await .unwrap(); let unit = validate_index_schema(reader.schema()).unwrap(); @@ -850,10 +905,11 @@ mod tests { assert_eq!(reader.row_groups_to_read(&pruning_predicate), vec![1]); // An empty time range yields no series without opening the file. + let (missing_index, _receiver) = index_handle(&metadata, &object_store); let empty = SeriesIndexSearcher::try_new( metadata, object_store, - "does-not-need-to-exist.parquet", + missing_index, None, Some(TimestampRange::empty()), ) diff --git a/src/mito2/src/series_index/version.rs b/src/mito2/src/series_index/version.rs index 4db883e7626..c31dfde906b 100644 --- a/src/mito2/src/series_index/version.rs +++ b/src/mito2/src/series_index/version.rs @@ -58,6 +58,11 @@ impl SeriesIndexFileHandle { } } + /// Returns the region and file identity used for storage and deletion. + pub(crate) fn file_id(&self) -> RegionFileId { + self.inner.file_id + } + pub(crate) fn entry(&self) -> &SeriesIndexEntry { &self.inner.entry } diff --git a/src/mito2/src/sst/parquet/file_range.rs b/src/mito2/src/sst/parquet/file_range.rs index f55aa524165..a3d4173304f 100644 --- a/src/mito2/src/sst/parquet/file_range.rs +++ b/src/mito2/src/sst/parquet/file_range.rs @@ -16,7 +16,7 @@ //! is usually a row group in a parquet file. use std::collections::HashMap; -use std::ops::BitAnd; +use std::ops::{BitAnd, Range}; use std::sync::Arc; use api::v1::{OpType, SemanticType}; @@ -29,19 +29,21 @@ use datatypes::arrow::record_batch::RecordBatch; use datatypes::schema::Schema; use futures::StreamExt; use mito_codec::row_converter::PrimaryKeyCodec; +use object_store::ObjectStore; use parquet::arrow::arrow_reader::RowSelection; use parquet::file::metadata::ParquetMetaData; use parquet::file::statistics::Statistics; -use snafu::{OptionExt, ResultExt}; +use snafu::{OptionExt, ResultExt, ensure}; use store_api::codec::PrimaryKeyEncoding; use store_api::metadata::RegionMetadataRef; use store_api::storage::{ColumnId, TimeSeriesRowSelector}; use table::predicate::Predicate; +use tokio::sync::OnceCell; use crate::cache::CacheStrategy; use crate::error::{ - ComputeArrowSnafu, DecodeStatsSnafu, EvalPartitionFilterSnafu, NewRecordBatchSnafu, - RecordBatchSnafu, Result, StatsNotPresentSnafu, UnexpectedSnafu, + ComputeArrowSnafu, DecodeStatsSnafu, EvalPartitionFilterSnafu, InvalidRecordBatchSnafu, + NewRecordBatchSnafu, RecordBatchSnafu, Result, StatsNotPresentSnafu, UnexpectedSnafu, }; use crate::read::compat::FlatCompatBatch; use crate::read::flat_projection::CompactionProjectionMapper; @@ -59,7 +61,9 @@ use crate::sst::parquet::reader::{ SimpleFilterContext, }; use crate::sst::parquet::row_group::ParquetFetchMetrics; +use crate::sst::parquet::row_selection::{intersect_row_selections, row_selection_from_row_ranges}; use crate::sst::parquet::stats::RowGroupPruningStats; +use crate::sst::range_index::{SstRangeIndexSearcher, range_index_path}; /// Checks if a row group contains delete operations by examining the min value of op_type column. /// @@ -100,6 +104,11 @@ pub struct FileRange { } impl FileRange { + /// Returns the shared range-index searcher, opening it on first use. + pub(crate) async fn range_index_searcher(&self) -> Result> { + self.context.range_index_searcher().await + } + /// Returns the region metadata stored in this SST. pub(crate) fn region_metadata(&self) -> &RegionMetadataRef { self.context.read_format().metadata() @@ -340,6 +349,40 @@ impl FileRange { Ok(Some(FlatRowGroupReader::new(self.context.clone(), stream))) } + /// Returns the source SST row-group index. + pub(crate) fn row_group_index(&self) -> usize { + self.row_group_idx + } + + /// Builds a reader from absolute row-group offsets supplied by a range index. + /// Existing pruning is intersected before reading, without predicate prefiltering. + pub(crate) async fn reader_by_row_ranges( + &self, + ranges: Vec>, + fetch_metrics: Option<&ParquetFetchMetrics>, + ) -> Result> { + let num_rows = self + .context + .reader_builder + .parquet_metadata() + .row_group(self.row_group_idx) + .num_rows() as usize; + let Some(selected) = refine_row_range_selection(ranges, num_rows, &self.row_selection)? + else { + return Ok(None); + }; + let stream = self + .context + .reader_builder + .build_without_prefilter(self.context.build_context( + self.row_group_idx, + Some(selected), + fetch_metrics, + )) + .await?; + Ok(Some(FlatRowGroupReader::new(self.context.clone(), stream))) + } + /// Returns the helper to compat batches. pub(crate) fn compat_batch(&self) -> Option<&FlatCompatBatch> { self.context.compat_batch() @@ -372,6 +415,32 @@ impl FileRange { } } +/// Intersects source-row offsets, rather than offsets within an already selected stream. +fn refine_row_range_selection( + ranges: Vec>, + num_rows: usize, + original: &Option, +) -> Result> { + let mut previous_end = 0; + for range in &ranges { + ensure!( + range.start >= previous_end && range.start < range.end && range.end <= num_rows, + InvalidRecordBatchSnafu { + reason: format!( + "invalid range-index row range {range:?} after {previous_end}, row group has {num_rows} rows" + ), + } + ); + previous_end = range.end; + } + let selected = row_selection_from_row_ranges(ranges.into_iter(), num_rows); + let selected = match original { + Some(original) => intersect_row_selections(original, &selected), + None => selected, + }; + Ok((selected.row_count() > 0).then_some(selected)) +} + fn refine_primary_key_selection( masks: &[BooleanArray], original: &Option, @@ -389,6 +458,10 @@ fn refine_primary_key_selection( /// Context shared by ranges of the same parquet SST. pub struct FileRangeContext { + /// Store for a range index registered in the scan's index snapshot. + range_index_store: Option, + /// Lazily opened range index shared by all ranges of this file. + range_index_searcher: OnceCell, /// Row group reader builder for the file. reader_builder: RowGroupReaderBuilder, /// Base of the context. @@ -399,13 +472,34 @@ pub type FileRangeContextRef = Arc; impl FileRangeContext { /// Creates a new [FileRangeContext]. - pub(crate) fn new(reader_builder: RowGroupReaderBuilder, base: RangeBase) -> Self { + pub(crate) fn new( + reader_builder: RowGroupReaderBuilder, + base: RangeBase, + range_index_store: Option, + ) -> Self { Self { reader_builder, base, + range_index_store, + range_index_searcher: OnceCell::new(), } } + /// Opens the range index once, retaining the SST handle throughout its use. + async fn range_index_searcher(&self) -> Result> { + let Some(store) = &self.range_index_store else { + return Ok(None); + }; + self.range_index_searcher + .get_or_try_init(|| async { + let file = self.reader_builder.file_handle(); + let path = range_index_path(file.region_id(), file.file_id().file_id()); + SstRangeIndexSearcher::open(store.clone(), &path).await + }) + .await + .map(Some) + } + /// Returns filters pushed down. pub(crate) fn filters(&self) -> &[SimpleFilterContext] { &self.base.filters @@ -969,6 +1063,56 @@ mod tests { ); } + #[test] + fn test_refine_row_range_selection() { + let original = Some(RowSelection::from(vec![ + RowSelector::skip(2), + RowSelector::select(3), + RowSelector::skip(1), + RowSelector::select(2), + ])); + let selected = refine_row_range_selection(vec![0..3, 4..7, 9..10], 10, &original) + .unwrap() + .unwrap(); + assert_eq!( + selected, + RowSelection::from(vec![ + RowSelector::skip(2), + RowSelector::select(1), + RowSelector::skip(1), + RowSelector::select(1), + RowSelector::skip(1), + RowSelector::select(1), + RowSelector::skip(1), + ]) + ); + assert!( + refine_row_range_selection(vec![0..2, 8..10], 10, &original) + .unwrap() + .is_none() + ); + assert!( + refine_row_range_selection(vec![], 10, &None) + .unwrap() + .is_none() + ); + assert_eq!( + refine_row_range_selection(vec![2..4, 4..6], 10, &None) + .unwrap() + .unwrap(), + RowSelection::from(vec![RowSelector::skip(2), RowSelector::select(4)]) + ); + for ranges in [ + std::iter::once(0..11).collect(), + std::iter::once(11..12).collect(), + std::iter::once(3..3).collect(), + vec![4..6, 5..7], + vec![4..6, 0..2], + ] { + assert!(refine_row_range_selection(ranges, 10, &None).is_err()); + } + } + #[test] fn test_refine_primary_key_selection_intersects_original_selection() { let original = Some(RowSelection::from(vec![ diff --git a/src/mito2/src/sst/parquet/reader.rs b/src/mito2/src/sst/parquet/reader.rs index e3232c189d6..f54812f97d4 100644 --- a/src/mito2/src/sst/parquet/reader.rs +++ b/src/mito2/src/sst/parquet/reader.rs @@ -68,6 +68,7 @@ use crate::metrics::{ use crate::read::flat_projection::CompactionProjectionMapper; use crate::read::prune::FlatPruneReader; use crate::read::read_columns::ReadColumns; +use crate::series_index::SeriesIndexReadContext; use crate::sst::file::FileHandle; use crate::sst::index::bloom_filter::applier::{ BloomFilterIndexApplierRef, BloomFilterIndexApplyMetrics, @@ -191,6 +192,8 @@ macro_rules! handle_index_error { /// Parquet SST reader builder. pub struct ParquetReaderBuilder { + /// Index snapshot and store used to determine range-index availability. + series_index: Option, /// SST directory. table_dir: String, /// Path type for generating file paths. @@ -247,6 +250,7 @@ impl ParquetReaderBuilder { object_store: ObjectStore, ) -> ParquetReaderBuilder { ParquetReaderBuilder { + series_index: None, table_dir, path_type, file_handle, @@ -274,6 +278,13 @@ impl ParquetReaderBuilder { } } + /// Sets the captured series-index context for range-index reads. + #[must_use] + pub(crate) fn series_index(mut self, context: Option) -> Self { + self.series_index = context; + self + } + /// Sets the scan-wide hint for rows in a decoded batch. #[must_use] pub(crate) fn batch_size(mut self, batch_size: usize) -> Self { @@ -692,6 +703,19 @@ impl ParquetReaderBuilder { let partition_filter = self.build_partition_filter(&read_format, &prune_schema)?; + let range_index_store = self.series_index.as_ref().and_then(|context| { + let region_id = self + .expected_metadata + .as_ref() + .unwrap_or(®ion_meta) + .region_id; + (self.file_handle.region_id() == region_id + && context + .version + .range_indexes + .contains(&self.file_handle.file_id().file_id())) + .then(|| context.store.clone()) + }); let context = FileRangeContext::new( reader_builder, RangeBase { @@ -706,6 +730,7 @@ impl ParquetReaderBuilder { pre_filter_mode: self.pre_filter_mode, partition_filter, }, + range_index_store, ); metrics.build_cost += start.elapsed(); diff --git a/src/mito2/src/sst/parquet/row_selection.rs b/src/mito2/src/sst/parquet/row_selection.rs index e31e4922a05..3ad40d8c74b 100644 --- a/src/mito2/src/sst/parquet/row_selection.rs +++ b/src/mito2/src/sst/parquet/row_selection.rs @@ -516,7 +516,7 @@ impl RowGroupSelection { /// /// returned: NNNNNNNNY (modified) /// NNNNNNNNYYNYN (original) -fn intersect_row_selections(left: &RowSelection, right: &RowSelection) -> RowSelection { +pub(crate) fn intersect_row_selections(left: &RowSelection, right: &RowSelection) -> RowSelection { let mut l_iter = left.iter().copied().peekable(); let mut r_iter = right.iter().copied().peekable();