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 <realevenyag@gmail.com>

* perf: use range indexes in two-phase series reader

Signed-off-by: evenyag <realevenyag@gmail.com>

* fix(mito2): update series index test fixtures after rebase

Signed-off-by: evenyag <realevenyag@gmail.com>

* refactor(mito2): move lazy range index searchers into file contexts

Signed-off-by: evenyag <realevenyag@gmail.com>

* fix: pin series index snapshot before data snapshot

Signed-off-by: evenyag <realevenyag@gmail.com>

* fix: preserve builder caching for index-covered SSTs

Signed-off-by: evenyag <realevenyag@gmail.com>

---------

Signed-off-by: evenyag <realevenyag@gmail.com>
This commit is contained in:
Yingwen
2026-09-17 08:45:00 +00:00
committed by GitHub
parent 39e34c0b74
commit 26c3281f2b
15 changed files with 1038 additions and 91 deletions
+10
View File
@@ -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())
+170 -2
View File
@@ -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::<Vec<_>>()
.await
.unwrap();
assert_eq!(
1,
batches.iter().map(|batch| batch.num_rows()).sum::<usize>()
);
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::<Vec<_>>().await.is_err());
}
}
}
/// Scans all partitions in round-robin fashion and returns rows sorted by (tag, ts).
+7 -1
View File
@@ -87,6 +87,12 @@ impl PartitionPruner {
}
}
/// Excludes files replaced by another candidate source from prefetching.
pub(crate) fn excluding_files(mut self, excluded: &HashSet<usize>) -> 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()
+29
View File
@@ -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<SeriesIndexReadContext>,
/// 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<SeriesIndexReadContext>,
) -> 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<SeriesIndexReadContext>,
/// 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<SeriesIndexReadContext>,
) -> 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<TimestampRange>) -> 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())
+423 -7
View File
@@ -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<StreamContext>,
partitions: Vec<Vec<PartitionRange>>,
partition_pruner: Arc<PartitionPruner>,
candidate_pruner: Arc<PartitionPruner>,
coverage: Arc<SeriesIndexCoverage>,
range_semaphore: Arc<Semaphore>,
memory_pool: Arc<dyn MemoryPool>,
metrics_set: ExecutionPlanMetricsSet,
@@ -108,11 +115,21 @@ impl SeriesCandidateScanner {
);
let all_ranges = partitions.iter().flatten().copied().collect::<Vec<_>>();
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::<Vec<_>>();
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<SeriesIndexFileHandle>,
/// Indices into `ScanInput.files`, shared by every occurrence of an SST.
covered_files: HashSet<usize>,
}
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<StreamContext>,
context: SeriesIndexReadContext,
index: SeriesIndexFileHandle,
semaphore: Arc<Semaphore>,
) -> 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<StreamContext>,
coverage: Arc<SeriesIndexCoverage>,
partition_pruner: Arc<PartitionPruner>,
range_semaphore: Arc<Semaphore>,
memory_pool: Arc<dyn MemoryPool>,
@@ -196,7 +359,18 @@ impl SeriesCandidateRangeBuilder {
part_range: PartitionRange,
merge_partition: usize,
) -> Result<BoxedRecordBatchStream> {
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<Pruner>) {
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::<Vec<_>>()
.await
.unwrap();
assert_eq!(
(0..1001)
.map(|tsid| MetricSeriesId { table_id: 1, tsid })
.collect::<Vec<_>>(),
groups.into_iter().flatten().collect::<Vec<_>>()
);
// 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::<Vec<_>>().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::<Vec<_>>()
.await
.is_err()
);
}
#[tokio::test]
async fn candidate_scanner_rejects_predicate_prefilter_pruner() {
let env = SchedulerEnv::new().await;
+19 -15
View File
@@ -159,15 +159,20 @@ impl SeriesBatchCollector {
struct MetricSeriesFilter {
range: SeriesRange,
series: Arc<HashSet<MetricSeriesId>>,
sorted_series: Arc<Vec<MetricSeriesId>>,
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<StreamContext>,
partition_ranges: Vec<PartitionRange>,
range: SeriesRange,
filter: MetricSeriesFilter,
codec: SparsePrimaryKeyCodec,
partition_pruner: Arc<PartitionPruner>,
@@ -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<StreamContext>,
part_range: PartitionRange,
range: SeriesRange,
filter: MetricSeriesFilter,
codec: SparsePrimaryKeyCodec,
partition_pruner: Arc<PartitionPruner>,
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
+5
View File
@@ -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<ObjectStore>,
/// 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(),
+2
View File
@@ -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(
+18 -2
View File
@@ -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<SeriesIndexVersion>,
}
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";
+61
View File
@@ -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<i64, WindowSequence>,
}
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<SeriesIndexEntry>,
@@ -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();
+114 -58
View File
@@ -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<ParquetIndexReader>,
@@ -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<TimestampRange>,
) -> Result<Self> {
@@ -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<PurgeRequest>) {
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<u8>],
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<PurgeRequest>) {
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<PurgeRequest>) {
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<MetricSeriesId> {
@@ -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::<Vec<_>>();
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::<Vec<_>>();
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()),
)
+5
View File
@@ -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
}
+149 -5
View File
@@ -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<Option<&SstRangeIndexSearcher>> {
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<Range<usize>>,
fetch_metrics: Option<&ParquetFetchMetrics>,
) -> Result<Option<FlatRowGroupReader>> {
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<Range<usize>>,
num_rows: usize,
original: &Option<RowSelection>,
) -> Result<Option<RowSelection>> {
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<RowSelection>,
@@ -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<ObjectStore>,
/// Lazily opened range index shared by all ranges of this file.
range_index_searcher: OnceCell<SstRangeIndexSearcher>,
/// Row group reader builder for the file.
reader_builder: RowGroupReaderBuilder,
/// Base of the context.
@@ -399,13 +472,34 @@ pub type FileRangeContextRef = Arc<FileRangeContext>;
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<ObjectStore>,
) -> 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<Option<&SstRangeIndexSearcher>> {
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![
+25
View File
@@ -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<SeriesIndexReadContext>,
/// 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<SeriesIndexReadContext>) -> 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(&region_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();
+1 -1
View File
@@ -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();