mirror of
https://github.com/GreptimeTeam/greptimedb.git
synced 2026-09-08 22:48:58 +00:00
feat(mito2): add range-based metric series reader (#8703)
* feat(mito2): add range-based metric series reader Signed-off-by: evenyag <realevenyag@gmail.com> * fix(mito2): address series reader review feedback Signed-off-by: evenyag <realevenyag@gmail.com> * fix(mito2): align series predicate filtering Signed-off-by: evenyag <realevenyag@gmail.com> * refactor: update semaphore usage and move prefilter flag Signed-off-by: evenyag <realevenyag@gmail.com> * fix(mito2): enforce candidate pruner invariant Signed-off-by: evenyag <realevenyag@gmail.com> --------- Signed-off-by: evenyag <realevenyag@gmail.com>
This commit is contained in:
@@ -419,6 +419,7 @@ fn apply_combined_filters(
|
||||
let predicate_mask = context.base.compute_filter_mask_flat(
|
||||
&record_batch,
|
||||
skip_fields,
|
||||
false,
|
||||
&mut tag_decode_state,
|
||||
)?;
|
||||
// If predicate filters out the entire batch, return None early
|
||||
|
||||
@@ -223,9 +223,9 @@ impl FlatPruneReader {
|
||||
}
|
||||
|
||||
let num_rows_before_filter = record_batch.num_rows();
|
||||
let Some(filtered_batch) = self
|
||||
.context
|
||||
.precise_filter_flat(record_batch, self.skip_fields)?
|
||||
let Some(filtered_batch) =
|
||||
self.context
|
||||
.precise_filter_flat(record_batch, self.skip_fields, false)?
|
||||
else {
|
||||
// the entire batch is filtered out
|
||||
self.metrics.filter_metrics.rows_precise_filtered += num_rows_before_filter;
|
||||
|
||||
@@ -178,6 +178,30 @@ impl PartitionPruner {
|
||||
}
|
||||
}
|
||||
|
||||
/// Options to create a [`Pruner`].
|
||||
#[derive(Debug, Clone, Copy)]
|
||||
pub struct PrunerOptions {
|
||||
/// Keeps file range builders until the query-scoped pruner is dropped.
|
||||
pub retain_builders: bool,
|
||||
/// Whether [`FileRange`]s execute the reduced-column predicate prefilter.
|
||||
///
|
||||
/// The flag is baked into the cached [`FileRangeBuilder`], so it belongs to
|
||||
/// the whole scan. All partitions sharing this pruner must agree on it.
|
||||
/// Two-stage series scans disable this prefilter: candidate discovery applies
|
||||
/// tag predicates, then the data phase calls [`FileRange::precise_filter_flat`]
|
||||
/// with tag filtering skipped so the predicates are applied exactly once.
|
||||
pub enable_predicate_prefilter: bool,
|
||||
}
|
||||
|
||||
impl Default for PrunerOptions {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
retain_builders: false,
|
||||
enable_predicate_prefilter: true,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// A pruner that prunes files for all partitions of a scanner.
|
||||
pub struct Pruner {
|
||||
/// Channels to send requests to workers.
|
||||
@@ -200,6 +224,8 @@ struct PrunerInner {
|
||||
manifest_pruned_files: Vec<AtomicBool>,
|
||||
/// Keeps file range builders until the query-scoped pruner is dropped.
|
||||
retain_builders: bool,
|
||||
/// Whether FileRanges should execute the reduced-column predicate prefilter.
|
||||
enable_predicate_prefilter: bool,
|
||||
}
|
||||
|
||||
impl Drop for PrunerInner {
|
||||
@@ -275,15 +301,19 @@ impl Pruner {
|
||||
/// Initially all file_entries have `remaining_ranges = 0`.
|
||||
/// Call `add_partition_ranges()` to initialize ref counts.
|
||||
pub fn new(stream_ctx: Arc<StreamContext>, num_workers: usize) -> Self {
|
||||
Self::new_with_retained_builders(stream_ctx, num_workers, false)
|
||||
Self::new_with_options(stream_ctx, num_workers, PrunerOptions::default())
|
||||
}
|
||||
|
||||
/// Creates a new pruner and optionally retains file range builders until drop.
|
||||
pub fn new_with_retained_builders(
|
||||
/// Creates a new pruner with the given options.
|
||||
pub fn new_with_options(
|
||||
stream_ctx: Arc<StreamContext>,
|
||||
num_workers: usize,
|
||||
retain_builders: bool,
|
||||
options: PrunerOptions,
|
||||
) -> Self {
|
||||
let PrunerOptions {
|
||||
retain_builders,
|
||||
enable_predicate_prefilter,
|
||||
} = options;
|
||||
let num_files = stream_ctx.input.num_files();
|
||||
let file_entries: Vec<_> = (0..num_files)
|
||||
.map(|_| {
|
||||
@@ -311,6 +341,7 @@ impl Pruner {
|
||||
stream_ctx,
|
||||
manifest_pruned_files,
|
||||
retain_builders,
|
||||
enable_predicate_prefilter,
|
||||
});
|
||||
|
||||
// Spawn worker tasks with their receivers
|
||||
@@ -486,6 +517,12 @@ impl Pruner {
|
||||
let _ = self.worker_senders[worker_idx].try_send(request);
|
||||
}
|
||||
|
||||
/// Returns whether this pruner builds file ranges with the reduced-column
|
||||
/// predicate prefilter enabled.
|
||||
pub fn predicate_prefilter_enabled(&self) -> bool {
|
||||
self.inner.enable_predicate_prefilter
|
||||
}
|
||||
|
||||
fn get_worker_idx(&self, file_id: FileId) -> usize {
|
||||
let file_id_hash = Uuid::from(file_id).as_u128() as usize;
|
||||
file_id_hash % self.inner.num_workers
|
||||
@@ -516,7 +553,13 @@ impl Pruner {
|
||||
.inner
|
||||
.stream_ctx
|
||||
.input
|
||||
.prune_file_after_manifest_check(file, pre_filter_mode, predicate, reader_metrics)
|
||||
.prune_file_after_manifest_check(
|
||||
file,
|
||||
pre_filter_mode,
|
||||
self.inner.enable_predicate_prefilter,
|
||||
predicate,
|
||||
reader_metrics,
|
||||
)
|
||||
.await?;
|
||||
|
||||
let arc_builder = Arc::new(builder);
|
||||
@@ -599,7 +642,13 @@ impl Pruner {
|
||||
inner
|
||||
.stream_ctx
|
||||
.input
|
||||
.prune_file_after_manifest_check(file, pre_filter_mode, predicate, &mut metrics)
|
||||
.prune_file_after_manifest_check(
|
||||
file,
|
||||
pre_filter_mode,
|
||||
inner.enable_predicate_prefilter,
|
||||
predicate,
|
||||
&mut metrics,
|
||||
)
|
||||
.await
|
||||
};
|
||||
|
||||
@@ -796,10 +845,13 @@ mod tests {
|
||||
.with_files(files)
|
||||
.with_append_mode(true);
|
||||
let stream_ctx = Arc::new(StreamContext::unordered_scan_ctx(input));
|
||||
let pruner = Arc::new(Pruner::new_with_retained_builders(
|
||||
let pruner = Arc::new(Pruner::new_with_options(
|
||||
stream_ctx,
|
||||
1,
|
||||
retain_builders,
|
||||
PrunerOptions {
|
||||
retain_builders,
|
||||
..Default::default()
|
||||
},
|
||||
));
|
||||
(env, pruner)
|
||||
}
|
||||
|
||||
@@ -501,7 +501,6 @@ pub(crate) fn build_candidate_range_cache_key(
|
||||
}
|
||||
|
||||
/// Builds a cache key for a two-phase series-data partition-range result.
|
||||
#[allow(dead_code)]
|
||||
pub(crate) fn build_series_range_cache_key(
|
||||
stream_ctx: &StreamContext,
|
||||
part_range: &PartitionRange,
|
||||
|
||||
@@ -1269,7 +1269,7 @@ impl ScanInput {
|
||||
return Ok(FileRangeBuilder::default());
|
||||
}
|
||||
|
||||
self.prune_file_after_manifest_check(file, pre_filter_mode, predicate, reader_metrics)
|
||||
self.prune_file_after_manifest_check(file, pre_filter_mode, true, predicate, reader_metrics)
|
||||
.await
|
||||
}
|
||||
|
||||
@@ -1284,6 +1284,7 @@ impl ScanInput {
|
||||
&self,
|
||||
file: &FileHandle,
|
||||
pre_filter_mode: PreFilterMode,
|
||||
enable_predicate_prefilter: bool,
|
||||
predicate: Option<Predicate>,
|
||||
reader_metrics: &mut ReaderMetrics,
|
||||
) -> Result<FileRangeBuilder> {
|
||||
@@ -1320,6 +1321,7 @@ impl ScanInput {
|
||||
.expected_metadata(Some(self.mapper.metadata().clone()))
|
||||
.compaction(self.compaction)
|
||||
.pre_filter_mode(pre_filter_mode)
|
||||
.enable_predicate_prefilter(enable_predicate_prefilter)
|
||||
.decode_primary_key_values(decode_pk_values)
|
||||
.build_reader_input(reader_metrics)
|
||||
.await;
|
||||
|
||||
@@ -88,6 +88,10 @@ impl SeriesCandidateScanner {
|
||||
///
|
||||
/// Callers must fall back to the legacy series-scan path when the scan contains
|
||||
/// extension ranges. Candidate-series discovery does not support other range types.
|
||||
///
|
||||
/// `pruner` must be built with `PrunerOptions::enable_predicate_prefilter` set to
|
||||
/// `false`: both scan phases share its file range builders, and the two-stage path
|
||||
/// applies simple filters on the precise-filter path instead.
|
||||
pub(crate) fn try_new(
|
||||
stream_ctx: Arc<StreamContext>,
|
||||
partitions: Vec<Vec<PartitionRange>>,
|
||||
@@ -106,6 +110,15 @@ impl SeriesCandidateScanner {
|
||||
reason: "candidate-series scan does not support extension ranges; use the legacy series-scan path",
|
||||
}
|
||||
);
|
||||
ensure!(
|
||||
!pruner.predicate_prefilter_enabled(),
|
||||
UnexpectedSnafu {
|
||||
reason: format!(
|
||||
"candidate-series scan for region {} requires a pruner without predicate prefiltering",
|
||||
stream_ctx.input.region_metadata().region_id
|
||||
),
|
||||
}
|
||||
);
|
||||
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));
|
||||
@@ -253,12 +266,15 @@ impl SeriesCandidateRangeBuilder {
|
||||
);
|
||||
let filter = build_primary_key_filter(
|
||||
&metadata,
|
||||
None,
|
||||
self.stream_ctx.input.predicate_group().predicate(),
|
||||
);
|
||||
return Ok(Some(candidate_primary_key_stream(Box::pin(raw), filter)));
|
||||
}
|
||||
|
||||
if self.stream_ctx.is_file_range_index(index) {
|
||||
let file = self.stream_ctx.input.file_from_index(index);
|
||||
let predicate = self.stream_ctx.input.predicate_for_file(file);
|
||||
if self
|
||||
.partition_pruner
|
||||
.try_skip_manifest_pruned_file_range(index, &self.part_metrics)
|
||||
@@ -277,9 +293,15 @@ impl SeriesCandidateRangeBuilder {
|
||||
self.part_metrics
|
||||
.merge_reader_metrics(&reader_metrics, None);
|
||||
|
||||
// Reuse the exact encoded-PK filter selected by the file's reader plan, but
|
||||
// execute it in `candidate_primary_key_stream` after the PK-only read.
|
||||
let filter = ranges.first().and_then(|range| range.primary_key_filter());
|
||||
// Build a fresh encoded-PK filter from this SST's metadata so schema
|
||||
// compatibility is evaluated independently for each file.
|
||||
let filter = ranges.first().and_then(|range| {
|
||||
build_primary_key_filter(
|
||||
range.region_metadata(),
|
||||
Some(metadata.as_ref()),
|
||||
predicate.as_ref(),
|
||||
)
|
||||
});
|
||||
let part_metrics = self.part_metrics.clone();
|
||||
let raw = Box::pin(try_stream! {
|
||||
let fetch_metrics = part_metrics
|
||||
@@ -547,6 +569,8 @@ fn decode_metric_series(
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::time::Instant;
|
||||
|
||||
use datafusion::execution::memory_pool::UnboundedMemoryPool;
|
||||
use datafusion_expr::{col, lit};
|
||||
use datatypes::arrow::array::{ArrayRef, DictionaryArray, UInt32Array};
|
||||
@@ -556,8 +580,55 @@ mod tests {
|
||||
use table::predicate::Predicate;
|
||||
|
||||
use super::*;
|
||||
use crate::read::flat_projection::FlatProjectionMapper;
|
||||
use crate::read::scan_region::ScanInput;
|
||||
use crate::read::scan_util::PartitionMetrics;
|
||||
use crate::test_util::scheduler_util::SchedulerEnv;
|
||||
use crate::test_util::sst_util::sst_region_metadata_with_encoding;
|
||||
|
||||
#[tokio::test]
|
||||
async fn candidate_scanner_rejects_predicate_prefilter_pruner() {
|
||||
let env = SchedulerEnv::new().await;
|
||||
let metadata = Arc::new(sst_region_metadata_with_encoding(
|
||||
PrimaryKeyEncoding::Sparse,
|
||||
));
|
||||
let mapper =
|
||||
FlatProjectionMapper::new(&metadata, 0..metadata.column_metadatas.len()).unwrap();
|
||||
let stream_ctx = Arc::new(StreamContext::seq_scan_ctx(ScanInput::new(
|
||||
env.access_layer.clone(),
|
||||
mapper,
|
||||
)));
|
||||
let pruner = Arc::new(Pruner::new(stream_ctx.clone(), 1));
|
||||
let metrics_set = ExecutionPlanMetricsSet::new();
|
||||
let part_metrics = PartitionMetrics::new(
|
||||
metadata.region_id,
|
||||
0,
|
||||
"candidate-test",
|
||||
Instant::now(),
|
||||
false,
|
||||
&metrics_set,
|
||||
);
|
||||
|
||||
let error = SeriesCandidateScanner::try_new(
|
||||
stream_ctx,
|
||||
Vec::new(),
|
||||
pruner,
|
||||
Arc::new(Semaphore::new(1)),
|
||||
Arc::new(UnboundedMemoryPool::default()),
|
||||
metrics_set,
|
||||
part_metrics,
|
||||
)
|
||||
.err()
|
||||
.unwrap();
|
||||
|
||||
assert!(matches!(error, crate::error::Error::Unexpected { .. }));
|
||||
assert!(
|
||||
error
|
||||
.to_string()
|
||||
.contains("requires a pruner without predicate prefiltering")
|
||||
);
|
||||
}
|
||||
|
||||
fn binary_batch(values: &[&[u8]]) -> RecordBatch {
|
||||
RecordBatch::try_new(
|
||||
primary_key_schema(),
|
||||
@@ -621,7 +692,7 @@ mod tests {
|
||||
let predicate = Predicate::new(vec![
|
||||
col(store_api::metric_engine_consts::DATA_SCHEMA_TABLE_ID_COLUMN_NAME).eq(lit(1_u32)),
|
||||
]);
|
||||
let filter = build_primary_key_filter(&metadata, Some(&predicate));
|
||||
let filter = build_primary_key_filter(&metadata, None, Some(&predicate));
|
||||
let input = Box::pin(futures::stream::iter(vec![Ok(dictionary_batch(
|
||||
&[table_1.as_slice(), table_2.as_slice()],
|
||||
&[0, 1],
|
||||
|
||||
@@ -12,9 +12,41 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
//! Shared types for range-based metric series reads.
|
||||
//! Reads selected metric series from partition ranges.
|
||||
|
||||
use std::collections::HashSet;
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
|
||||
use async_stream::try_stream;
|
||||
use futures::TryStreamExt;
|
||||
use mito_codec::row_converter::{PrimaryKeyFilter, SparsePrimaryKeyCodec};
|
||||
use snafu::ResultExt;
|
||||
use store_api::region_engine::PartitionRange;
|
||||
use tokio::sync::Semaphore;
|
||||
|
||||
#[cfg(feature = "enterprise")]
|
||||
use crate::error::InvalidRequestSnafu;
|
||||
use crate::error::{JoinSnafu, Result, UnexpectedSnafu};
|
||||
use crate::read::BoxedRecordBatchStream;
|
||||
use crate::read::pruner::PartitionPruner;
|
||||
use crate::read::range_cache::{
|
||||
build_series_range_cache_key, cache_flat_range_stream, cached_flat_range_stream,
|
||||
};
|
||||
use crate::read::scan_region::StreamContext;
|
||||
use crate::read::scan_util::{
|
||||
PartitionMetrics, SplitRecordBatchStream, compute_average_batch_size,
|
||||
compute_parallel_channel_size, new_filter_metrics, scan_flat_mem_ranges,
|
||||
should_split_flat_batches_for_merge,
|
||||
};
|
||||
use crate::read::seq_scan::SeqScan;
|
||||
use crate::read::series_candidate::{MetricSeriesId, validate_metric_metadata};
|
||||
use crate::sst::parquet::DEFAULT_READ_BATCH_SIZE;
|
||||
use crate::sst::parquet::flat_format::primary_key_column_index;
|
||||
use crate::sst::parquet::prefilter::prefilter_flat_batch_by_primary_key;
|
||||
use crate::sst::parquet::reader::ReaderMetrics;
|
||||
use crate::sst::parquet::row_group::ParquetFetchMetrics;
|
||||
|
||||
#[allow(dead_code)]
|
||||
const TSID_DOMAIN_END: u128 = 1u128 << u64::BITS;
|
||||
|
||||
/// A stable partition of the TSID integer domain.
|
||||
@@ -24,7 +56,6 @@ pub(crate) struct SeriesRange {
|
||||
end: u128,
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
impl SeriesRange {
|
||||
pub(crate) fn new(partition: usize, partitions: usize) -> Option<Self> {
|
||||
if partitions == 0 || partition >= partitions {
|
||||
@@ -35,28 +66,684 @@ impl SeriesRange {
|
||||
let partition = partition as u128;
|
||||
let boundary = |partition: u128| {
|
||||
let numerator = partition * TSID_DOMAIN_END;
|
||||
numerator / partitions + u128::from(!numerator.is_multiple_of(partitions))
|
||||
numerator.div_ceil(partitions)
|
||||
};
|
||||
Some(Self {
|
||||
start: boundary(partition),
|
||||
end: boundary(partition + 1),
|
||||
})
|
||||
}
|
||||
|
||||
fn partition_for(tsid: u64, partitions: usize) -> usize {
|
||||
((tsid as u128 * partitions as u128) >> u64::BITS) as usize
|
||||
}
|
||||
}
|
||||
|
||||
/// All series assigned to one data-reader partition.
|
||||
#[derive(Debug)]
|
||||
pub(crate) struct AssignedSeriesBatch {
|
||||
range: SeriesRange,
|
||||
series: Vec<MetricSeriesId>,
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
impl AssignedSeriesBatch {
|
||||
fn new(range: SeriesRange, series: Vec<MetricSeriesId>) -> Self {
|
||||
Self { range, series }
|
||||
}
|
||||
|
||||
pub(crate) fn range(&self) -> SeriesRange {
|
||||
self.range
|
||||
}
|
||||
|
||||
pub(crate) fn series(&self) -> &[MetricSeriesId] {
|
||||
&self.series
|
||||
}
|
||||
}
|
||||
|
||||
/// Collects candidate batches and assigns every TSID by its integer range.
|
||||
#[allow(dead_code)]
|
||||
pub(crate) struct SeriesBatchCollector {
|
||||
assignments: Vec<Vec<MetricSeriesId>>,
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
impl SeriesBatchCollector {
|
||||
pub(crate) fn new(partitions: usize) -> Option<Self> {
|
||||
(partitions > 0).then(|| Self {
|
||||
assignments: (0..partitions).map(|_| Vec::new()).collect(),
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) fn push(&mut self, batch: Vec<MetricSeriesId>) {
|
||||
let partitions = self.assignments.len();
|
||||
for series in batch {
|
||||
let partition = SeriesRange::partition_for(series.tsid, partitions);
|
||||
self.assignments[partition].push(series);
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn finish(self) -> Vec<AssignedSeriesBatch> {
|
||||
let partitions = self.assignments.len();
|
||||
self.assignments
|
||||
.into_iter()
|
||||
.enumerate()
|
||||
.map(|(partition, series)| {
|
||||
AssignedSeriesBatch::new(SeriesRange::new(partition, partitions).unwrap(), series)
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
}
|
||||
|
||||
/// Immutable allow-list for one partition's metric series.
|
||||
#[derive(Clone, Debug)]
|
||||
struct MetricSeriesFilter {
|
||||
range: SeriesRange,
|
||||
series: Arc<HashSet<MetricSeriesId>>,
|
||||
}
|
||||
|
||||
impl MetricSeriesFilter {
|
||||
fn new(assigned: &AssignedSeriesBatch) -> Self {
|
||||
let series = assigned.series.iter().copied().collect();
|
||||
Self {
|
||||
range: assigned.range,
|
||||
series: Arc::new(series),
|
||||
}
|
||||
}
|
||||
|
||||
fn primary_key_filter(&self, codec: SparsePrimaryKeyCodec) -> Box<dyn PrimaryKeyFilter> {
|
||||
Box::new(MetricSeriesPrimaryKeyFilter {
|
||||
codec,
|
||||
series: self.series.clone(),
|
||||
last_primary_key: Vec::new(),
|
||||
last_match: None,
|
||||
})
|
||||
}
|
||||
|
||||
fn overlaps_encoded_bounds(
|
||||
&self,
|
||||
codec: &SparsePrimaryKeyCodec,
|
||||
encoded_min: &[u8],
|
||||
encoded_max: &[u8],
|
||||
) -> Option<bool> {
|
||||
let (min_table_id, min_tsid) = codec.decode_ids(encoded_min).ok()?;
|
||||
let (max_table_id, max_tsid) = codec.decode_ids(encoded_max).ok()?;
|
||||
let min = MetricSeriesId {
|
||||
table_id: min_table_id,
|
||||
tsid: min_tsid,
|
||||
};
|
||||
let max = MetricSeriesId {
|
||||
table_id: max_table_id,
|
||||
tsid: max_tsid,
|
||||
};
|
||||
if min > max {
|
||||
return None;
|
||||
}
|
||||
if min_table_id != max_table_id {
|
||||
return Some(true);
|
||||
}
|
||||
|
||||
Some(u128::from(min_tsid) < self.range.end && u128::from(max_tsid) >= self.range.start)
|
||||
}
|
||||
}
|
||||
|
||||
struct MetricSeriesPrimaryKeyFilter {
|
||||
codec: SparsePrimaryKeyCodec,
|
||||
series: Arc<HashSet<MetricSeriesId>>,
|
||||
last_primary_key: Vec<u8>,
|
||||
last_match: Option<bool>,
|
||||
}
|
||||
|
||||
impl PrimaryKeyFilter for MetricSeriesPrimaryKeyFilter {
|
||||
fn matches(&mut self, primary_key: &[u8]) -> mito_codec::error::Result<bool> {
|
||||
if let Some(last_match) = self.last_match
|
||||
&& self.last_primary_key == primary_key
|
||||
{
|
||||
return Ok(last_match);
|
||||
}
|
||||
|
||||
let (table_id, tsid) = self.codec.decode_ids(primary_key)?;
|
||||
let matched = self.series.contains(&MetricSeriesId { table_id, tsid });
|
||||
self.last_primary_key.clear();
|
||||
self.last_primary_key.extend_from_slice(primary_key);
|
||||
self.last_match = Some(matched);
|
||||
Ok(matched)
|
||||
}
|
||||
}
|
||||
|
||||
fn filter_flat_stream_by_series(
|
||||
mut input: BoxedRecordBatchStream,
|
||||
codec: SparsePrimaryKeyCodec,
|
||||
filter: MetricSeriesFilter,
|
||||
) -> BoxedRecordBatchStream {
|
||||
Box::pin(try_stream! {
|
||||
let mut primary_key_filter = filter.primary_key_filter(codec);
|
||||
while let Some(batch) = input.try_next().await? {
|
||||
let pk_idx = primary_key_column_index(batch.num_columns());
|
||||
if let Some(batch) = prefilter_flat_batch_by_primary_key(
|
||||
batch,
|
||||
pk_idx,
|
||||
primary_key_filter.as_mut(),
|
||||
)? {
|
||||
yield batch;
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
/// Reads all collected metric series assigned to one partition.
|
||||
#[allow(dead_code)]
|
||||
pub(crate) struct SeriesReader {
|
||||
stream_ctx: Arc<StreamContext>,
|
||||
partition_ranges: Vec<PartitionRange>,
|
||||
range: SeriesRange,
|
||||
filter: MetricSeriesFilter,
|
||||
codec: SparsePrimaryKeyCodec,
|
||||
partition_pruner: Arc<PartitionPruner>,
|
||||
range_semaphore: Arc<Semaphore>,
|
||||
part_metrics: PartitionMetrics,
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
impl SeriesReader {
|
||||
/// Creates a reader for the series assigned to one data partition.
|
||||
///
|
||||
/// `partition_pruner` must come from the candidate scanner's pruner, which
|
||||
/// has predicate prefiltering disabled. The file read path applies precise
|
||||
/// filters with tag filtering skipped because candidate discovery has already
|
||||
/// enforced tag predicates through the exact assigned-series set.
|
||||
///
|
||||
/// A single `range_semaphore` covers both the range-build phase and the final
|
||||
/// merge.
|
||||
pub(crate) fn try_new(
|
||||
stream_ctx: Arc<StreamContext>,
|
||||
partition_ranges: Vec<PartitionRange>,
|
||||
assigned_series: AssignedSeriesBatch,
|
||||
partition_pruner: Arc<PartitionPruner>,
|
||||
range_semaphore: Arc<Semaphore>,
|
||||
part_metrics: PartitionMetrics,
|
||||
) -> Result<Self> {
|
||||
validate_metric_metadata(&stream_ctx)?;
|
||||
#[cfg(feature = "enterprise")]
|
||||
snafu::ensure!(
|
||||
stream_ctx.input.extension_ranges().is_empty(),
|
||||
InvalidRequestSnafu {
|
||||
region_id: stream_ctx.input.region_metadata().region_id,
|
||||
reason: "series reader does not support extension ranges",
|
||||
}
|
||||
);
|
||||
|
||||
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,
|
||||
range_semaphore,
|
||||
part_metrics,
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) async fn build_stream(&self) -> Result<BoxedRecordBatchStream> {
|
||||
if self.partition_ranges.is_empty() || self.filter.series.is_empty() {
|
||||
return Ok(Box::pin(futures::stream::empty()));
|
||||
}
|
||||
|
||||
let mut tasks = Vec::with_capacity(self.partition_ranges.len());
|
||||
for part_range in self.partition_ranges.iter().copied() {
|
||||
let stream_ctx = self.stream_ctx.clone();
|
||||
let filter = self.filter.clone();
|
||||
let codec = self.codec.clone();
|
||||
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 {
|
||||
reason: format!("failed to acquire series range permit: {error}"),
|
||||
}
|
||||
.build()
|
||||
})?;
|
||||
build_series_partition_range(
|
||||
stream_ctx,
|
||||
part_range,
|
||||
range,
|
||||
filter,
|
||||
codec,
|
||||
partition_pruner,
|
||||
part_metrics,
|
||||
)
|
||||
.await
|
||||
}));
|
||||
}
|
||||
|
||||
let mut range_streams = Vec::with_capacity(tasks.len());
|
||||
let mut estimated_batch_sizes = Vec::with_capacity(tasks.len());
|
||||
for task in tasks {
|
||||
let (stream, estimated_batch_size) = task.await.context(JoinSnafu)??;
|
||||
range_streams.push(stream);
|
||||
estimated_batch_sizes.push(estimated_batch_size);
|
||||
}
|
||||
|
||||
// Every range task above has finished, so all build permits are released
|
||||
// and the final merge can reuse the same semaphore.
|
||||
let estimated_batch_size = compute_average_batch_size(estimated_batch_sizes);
|
||||
SeqScan::build_flat_reader_from_sources(
|
||||
&self.stream_ctx,
|
||||
range_streams,
|
||||
Some(self.range_semaphore.clone()),
|
||||
Some(&self.part_metrics),
|
||||
true,
|
||||
compute_parallel_channel_size(estimated_batch_size),
|
||||
)
|
||||
.await
|
||||
}
|
||||
}
|
||||
|
||||
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 cache_key = build_series_range_cache_key(&stream_ctx, &part_range, range);
|
||||
if let Some(key) = cache_key.as_ref() {
|
||||
if let Some(value) = stream_ctx.input.cache_strategy.get_range_result(key) {
|
||||
part_metrics.inc_range_cache_hit();
|
||||
return Ok((cached_flat_range_stream(value), DEFAULT_READ_BATCH_SIZE));
|
||||
}
|
||||
part_metrics.inc_range_cache_miss();
|
||||
}
|
||||
|
||||
let range_meta = &stream_ctx.ranges[part_range.identifier];
|
||||
let split_batch_size = should_split_flat_batches_for_merge(&stream_ctx, range_meta);
|
||||
let mut sources = Vec::with_capacity(range_meta.row_group_indices.len());
|
||||
|
||||
for index in range_meta.row_group_indices.iter().copied() {
|
||||
if stream_ctx.is_mem_range_index(index) {
|
||||
let stream = Box::pin(scan_flat_mem_ranges(
|
||||
stream_ctx.clone(),
|
||||
part_metrics.clone(),
|
||||
index,
|
||||
range_meta.time_range,
|
||||
));
|
||||
sources.push(filter_flat_stream_by_series(
|
||||
stream,
|
||||
codec.clone(),
|
||||
filter.clone(),
|
||||
));
|
||||
continue;
|
||||
}
|
||||
|
||||
if stream_ctx.is_file_range_index(index) {
|
||||
let file = stream_ctx.input.file_from_index(index);
|
||||
if matches!(
|
||||
file.primary_key_range()
|
||||
.and_then(|(min, max)| { filter.overlaps_encoded_bounds(&codec, &min, &max) }),
|
||||
Some(false)
|
||||
) {
|
||||
continue;
|
||||
}
|
||||
|
||||
if partition_pruner.try_skip_manifest_pruned_file_range(index, &part_metrics) {
|
||||
continue;
|
||||
}
|
||||
let mut reader_metrics = ReaderMetrics {
|
||||
filter_metrics: new_filter_metrics(part_metrics.explain_verbose()),
|
||||
..Default::default()
|
||||
};
|
||||
let file_ranges = partition_pruner
|
||||
.build_file_ranges(index, &part_metrics, &mut reader_metrics)
|
||||
.await?;
|
||||
part_metrics.inc_num_file_ranges(file_ranges.len());
|
||||
part_metrics.merge_reader_metrics(&reader_metrics, None);
|
||||
let ranges = file_ranges
|
||||
.iter()
|
||||
.filter(|file_range| {
|
||||
!matches!(
|
||||
file_range.primary_key_range().and_then(|(min, max)| {
|
||||
filter.overlaps_encoded_bounds(&codec, min, max)
|
||||
}),
|
||||
Some(false)
|
||||
)
|
||||
})
|
||||
.cloned()
|
||||
.collect::<smallvec::SmallVec<[_; 2]>>();
|
||||
if ranges.is_empty() {
|
||||
continue;
|
||||
}
|
||||
|
||||
let stream = scan_series_file_ranges(
|
||||
part_metrics.clone(),
|
||||
ranges,
|
||||
filter.clone(),
|
||||
codec.clone(),
|
||||
);
|
||||
sources.push(Box::pin(stream) as BoxedRecordBatchStream);
|
||||
continue;
|
||||
}
|
||||
|
||||
return UnexpectedSnafu {
|
||||
reason: format!(
|
||||
"series reader received unsupported range index {}",
|
||||
index.index
|
||||
),
|
||||
}
|
||||
.fail();
|
||||
}
|
||||
|
||||
if split_batch_size.is_some() {
|
||||
sources = sources
|
||||
.into_iter()
|
||||
.map(|stream| Box::pin(SplitRecordBatchStream::new(stream)) as BoxedRecordBatchStream)
|
||||
.collect();
|
||||
}
|
||||
let estimated_batch_size = split_batch_size.unwrap_or(DEFAULT_READ_BATCH_SIZE);
|
||||
let stream = SeqScan::build_flat_reader_from_sources(
|
||||
&stream_ctx,
|
||||
sources,
|
||||
None,
|
||||
Some(&part_metrics),
|
||||
false,
|
||||
compute_parallel_channel_size(estimated_batch_size),
|
||||
)
|
||||
.await?;
|
||||
let stream = match cache_key {
|
||||
Some(key) => cache_flat_range_stream(
|
||||
stream,
|
||||
stream_ctx.input.cache_strategy.clone(),
|
||||
key,
|
||||
part_metrics,
|
||||
),
|
||||
None => stream,
|
||||
};
|
||||
Ok((stream, estimated_batch_size))
|
||||
}
|
||||
|
||||
fn scan_series_file_ranges(
|
||||
part_metrics: PartitionMetrics,
|
||||
ranges: smallvec::SmallVec<[crate::sst::parquet::file_range::FileRange; 2]>,
|
||||
filter: MetricSeriesFilter,
|
||||
codec: SparsePrimaryKeyCodec,
|
||||
) -> impl futures::Stream<Item = Result<datatypes::arrow::record_batch::RecordBatch>> {
|
||||
try_stream! {
|
||||
let fetch_metrics = part_metrics
|
||||
.explain_verbose()
|
||||
.then(|| Arc::new(ParquetFetchMetrics::default()));
|
||||
let mut reader_metrics = ReaderMetrics {
|
||||
fetch_metrics: fetch_metrics.clone(),
|
||||
..Default::default()
|
||||
};
|
||||
let mut primary_key_filter = filter.primary_key_filter(codec);
|
||||
|
||||
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 build_cost = build_start.elapsed();
|
||||
reader_metrics.build_cost += build_cost;
|
||||
part_metrics.inc_build_reader_cost(build_cost);
|
||||
|
||||
let scan_start = Instant::now();
|
||||
while let Some(record_batch) = reader.next_batch().await? {
|
||||
reader_metrics.num_record_batches += 1;
|
||||
reader_metrics.num_batches += 1;
|
||||
reader_metrics.num_rows += record_batch.num_rows();
|
||||
|
||||
let num_rows_before_filter = record_batch.num_rows();
|
||||
let Some(record_batch) = range.precise_filter_flat(
|
||||
record_batch,
|
||||
range.pre_filter_mode().skip_fields(),
|
||||
true,
|
||||
)? else {
|
||||
reader_metrics.filter_metrics.rows_precise_filtered +=
|
||||
num_rows_before_filter;
|
||||
continue;
|
||||
};
|
||||
reader_metrics.filter_metrics.rows_precise_filtered +=
|
||||
num_rows_before_filter - record_batch.num_rows();
|
||||
|
||||
let record_batch = if let Some(mapper) = range.compaction_projection_mapper() {
|
||||
mapper.project(record_batch)?
|
||||
} else {
|
||||
record_batch
|
||||
};
|
||||
if let Some(compat) = range.compat_batch() {
|
||||
yield compat.compat(record_batch)?;
|
||||
} else {
|
||||
yield record_batch;
|
||||
}
|
||||
}
|
||||
reader_metrics.scan_cost += scan_start.elapsed();
|
||||
}
|
||||
|
||||
reader_metrics.observe_rows("series_data");
|
||||
reader_metrics.filter_metrics.observe();
|
||||
part_metrics.merge_reader_metrics(&reader_metrics, None);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use store_api::codec::PrimaryKeyEncoding;
|
||||
|
||||
use super::*;
|
||||
use crate::error::DecodeSnafu;
|
||||
use crate::test_util::sst_util::sst_region_metadata_with_encoding;
|
||||
|
||||
fn series(table_id: u32, tsid: u64) -> MetricSeriesId {
|
||||
MetricSeriesId { table_id, tsid }
|
||||
}
|
||||
|
||||
fn assigned_batch(
|
||||
partitions: usize,
|
||||
partition: usize,
|
||||
series: Vec<MetricSeriesId>,
|
||||
) -> AssignedSeriesBatch {
|
||||
let mut collector = SeriesBatchCollector::new(partitions).unwrap();
|
||||
collector.push(series);
|
||||
collector.finish().remove(partition)
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn series_range_assignment_is_stable_across_batch_boundaries() {
|
||||
let input = vec![
|
||||
series(1, 0),
|
||||
series(2, 1u64 << 62),
|
||||
series(1, 1u64 << 63),
|
||||
series(2, 3u64 << 62),
|
||||
series(1, u64::MAX),
|
||||
series(3, 0),
|
||||
];
|
||||
let mut first = SeriesBatchCollector::new(4).unwrap();
|
||||
first.push(input.clone());
|
||||
let first = first.finish();
|
||||
|
||||
let mut second = SeriesBatchCollector::new(4).unwrap();
|
||||
for chunk in input.chunks(2) {
|
||||
second.push(chunk.to_vec());
|
||||
}
|
||||
let second = second.finish();
|
||||
|
||||
assert_eq!(
|
||||
first
|
||||
.iter()
|
||||
.map(AssignedSeriesBatch::series)
|
||||
.collect::<Vec<_>>(),
|
||||
second
|
||||
.iter()
|
||||
.map(AssignedSeriesBatch::series)
|
||||
.collect::<Vec<_>>()
|
||||
);
|
||||
assert_eq!(&[series(1, 0), series(3, 0)], first[0].series());
|
||||
assert_eq!(&[series(2, 1u64 << 62)], first[1].series());
|
||||
assert_eq!(
|
||||
input.len(),
|
||||
first
|
||||
.iter()
|
||||
.map(|batch| batch.series().len())
|
||||
.sum::<usize>()
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn series_ranges_cover_the_tsid_domain() {
|
||||
let ranges = (0..3)
|
||||
.map(|partition| SeriesRange::new(partition, 3).unwrap())
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
assert_eq!(0, ranges[0].start);
|
||||
assert_eq!(TSID_DOMAIN_END, ranges[2].end);
|
||||
assert_eq!(ranges[0].end, ranges[1].start);
|
||||
assert_eq!(ranges[1].end, ranges[2].start);
|
||||
assert_eq!(0, SeriesRange::partition_for(0, 3));
|
||||
assert_eq!(2, SeriesRange::partition_for(u64::MAX, 3));
|
||||
|
||||
for (partition, range) in ranges.iter().enumerate() {
|
||||
assert_eq!(partition, SeriesRange::partition_for(range.start as u64, 3));
|
||||
if range.end < TSID_DOMAIN_END {
|
||||
assert_eq!(
|
||||
partition + 1,
|
||||
SeriesRange::partition_for(range.end as u64, 3)
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn metric_series_filter_matches_encoded_primary_key() {
|
||||
let metadata = Arc::new(sst_region_metadata_with_encoding(
|
||||
PrimaryKeyEncoding::Sparse,
|
||||
));
|
||||
let assigned = assigned_batch(1, 0, vec![series(1, 10), series(2, 20)]);
|
||||
let filter = MetricSeriesFilter::new(&assigned);
|
||||
let codec = SparsePrimaryKeyCodec::new(&metadata);
|
||||
let mut primary_key_filter = filter.primary_key_filter(codec);
|
||||
|
||||
for selected in [series(1, 10), series(2, 20)] {
|
||||
assert!(
|
||||
primary_key_filter
|
||||
.matches(&encode_series(&metadata, selected.table_id, selected.tsid))
|
||||
.context(DecodeSnafu)
|
||||
.unwrap()
|
||||
);
|
||||
}
|
||||
for unselected in [series(2, 10), series(1, 20)] {
|
||||
assert!(
|
||||
!primary_key_filter
|
||||
.matches(&encode_series(
|
||||
&metadata,
|
||||
unselected.table_id,
|
||||
unselected.tsid,
|
||||
))
|
||||
.context(DecodeSnafu)
|
||||
.unwrap()
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
fn encode_series(
|
||||
metadata: &store_api::metadata::RegionMetadataRef,
|
||||
table_id: u32,
|
||||
tsid: u64,
|
||||
) -> Vec<u8> {
|
||||
let codec = SparsePrimaryKeyCodec::new(metadata);
|
||||
let mut primary_key = Vec::new();
|
||||
codec
|
||||
.encode_internal(table_id, tsid, &mut primary_key)
|
||||
.unwrap();
|
||||
primary_key
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn series_range_overlaps_single_table_primary_key_bounds() {
|
||||
let metadata = Arc::new(sst_region_metadata_with_encoding(
|
||||
PrimaryKeyEncoding::Sparse,
|
||||
));
|
||||
let assigned = assigned_batch(2, 0, vec![series(1, 10), series(2, 20)]);
|
||||
let filter = MetricSeriesFilter::new(&assigned);
|
||||
let codec = SparsePrimaryKeyCodec::new(&metadata);
|
||||
|
||||
// The full assigned range overlaps even though the actual candidate bounds do not.
|
||||
assert_eq!(
|
||||
Some(true),
|
||||
filter.overlaps_encoded_bounds(
|
||||
&codec,
|
||||
&encode_series(&metadata, 1, 100),
|
||||
&encode_series(&metadata, 1, 200),
|
||||
)
|
||||
);
|
||||
assert_eq!(
|
||||
Some(false),
|
||||
filter.overlaps_encoded_bounds(
|
||||
&codec,
|
||||
&encode_series(&metadata, 1, 1u64 << 63),
|
||||
&encode_series(&metadata, 1, u64::MAX),
|
||||
)
|
||||
);
|
||||
assert_eq!(
|
||||
Some(true),
|
||||
filter.overlaps_encoded_bounds(
|
||||
&codec,
|
||||
&encode_series(&metadata, 3, 10),
|
||||
&encode_series(&metadata, 3, 20),
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn series_range_keeps_multiple_table_primary_key_bounds() {
|
||||
let metadata = Arc::new(sst_region_metadata_with_encoding(
|
||||
PrimaryKeyEncoding::Sparse,
|
||||
));
|
||||
let assigned = assigned_batch(2, 0, vec![series(1, 10), series(2, 20)]);
|
||||
let filter = MetricSeriesFilter::new(&assigned);
|
||||
let codec = SparsePrimaryKeyCodec::new(&metadata);
|
||||
|
||||
assert_eq!(
|
||||
Some(true),
|
||||
filter.overlaps_encoded_bounds(
|
||||
&codec,
|
||||
&encode_series(&metadata, 1, 1u64 << 63),
|
||||
&encode_series(&metadata, 2, u64::MAX),
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn series_statistics_keep_invalid_or_inverted_bounds() {
|
||||
let metadata = Arc::new(sst_region_metadata_with_encoding(
|
||||
PrimaryKeyEncoding::Sparse,
|
||||
));
|
||||
let assigned = assigned_batch(2, 0, vec![series(1, 10)]);
|
||||
let filter = MetricSeriesFilter::new(&assigned);
|
||||
let codec = SparsePrimaryKeyCodec::new(&metadata);
|
||||
|
||||
assert_eq!(
|
||||
None,
|
||||
filter.overlaps_encoded_bounds(&codec, b"invalid", b"bounds")
|
||||
);
|
||||
assert_eq!(
|
||||
None,
|
||||
filter.overlaps_encoded_bounds(
|
||||
&codec,
|
||||
&encode_series(&metadata, 2, 0),
|
||||
&encode_series(&metadata, 1, 0),
|
||||
)
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -53,7 +53,7 @@ use crate::sst::parquet::flat_format::{
|
||||
time_index_column_index,
|
||||
};
|
||||
use crate::sst::parquet::json_align::ProjectedRecordBatchStream;
|
||||
use crate::sst::parquet::prefilter::{CachedPrimaryKeyFilter, primary_key_filter_mask};
|
||||
use crate::sst::parquet::prefilter::primary_key_filter_mask;
|
||||
use crate::sst::parquet::reader::{
|
||||
FlatRowGroupReader, MaybeFilter, RowGroupBuildContext, RowGroupReaderBuilder,
|
||||
SimpleFilterContext,
|
||||
@@ -100,13 +100,12 @@ pub struct FileRange {
|
||||
}
|
||||
|
||||
impl FileRange {
|
||||
/// Builds the encoded-primary-key filter selected for this file.
|
||||
pub(crate) fn primary_key_filter(&self) -> Option<CachedPrimaryKeyFilter> {
|
||||
self.context.reader_builder.primary_key_filter()
|
||||
/// Returns the region metadata stored in this SST.
|
||||
pub(crate) fn region_metadata(&self) -> &RegionMetadataRef {
|
||||
self.context.read_format().metadata()
|
||||
}
|
||||
|
||||
/// Returns encoded primary-key min/max statistics for this row group.
|
||||
#[allow(dead_code)]
|
||||
pub(crate) fn primary_key_range(&self) -> Option<(&[u8], &[u8])> {
|
||||
let metadata = self.context.reader_builder.parquet_metadata();
|
||||
let num_columns = metadata.file_metadata().schema_descr().num_columns();
|
||||
@@ -301,9 +300,8 @@ impl FileRange {
|
||||
/// filter and this range's existing row selection.
|
||||
///
|
||||
/// This deliberately bypasses generic predicate prefiltering. The series
|
||||
/// pruner selected the row group independently, and non-tag simple predicates
|
||||
/// are applied after merge and dedup by the series reader.
|
||||
#[allow(dead_code)]
|
||||
/// pruner selected the row group independently, and simple predicates retained
|
||||
/// by the disabled prefilter plan are applied precisely before merge.
|
||||
pub(crate) async fn reader_by_primary_key(
|
||||
&self,
|
||||
primary_key_filter: &mut dyn mito_codec::row_converter::PrimaryKeyFilter,
|
||||
@@ -348,13 +346,28 @@ impl FileRange {
|
||||
self.context.compaction_projection_mapper()
|
||||
}
|
||||
|
||||
/// Filters a full-projection batch using this range's precise filters.
|
||||
pub(crate) fn precise_filter_flat(
|
||||
&self,
|
||||
input: RecordBatch,
|
||||
skip_fields: bool,
|
||||
skip_tags: bool,
|
||||
) -> Result<Option<RecordBatch>> {
|
||||
self.context
|
||||
.precise_filter_flat(input, skip_fields, skip_tags)
|
||||
}
|
||||
|
||||
/// Returns the precise-filter mode configured for this range.
|
||||
pub(crate) fn pre_filter_mode(&self) -> PreFilterMode {
|
||||
self.context.pre_filter_mode()
|
||||
}
|
||||
|
||||
/// Returns the file handle of the file range.
|
||||
pub(crate) fn file_handle(&self) -> &FileHandle {
|
||||
self.context.reader_builder.file_handle()
|
||||
}
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
fn refine_primary_key_selection(
|
||||
masks: &[BooleanArray],
|
||||
original: &Option<RowSelection>,
|
||||
@@ -431,8 +444,9 @@ impl FileRangeContext {
|
||||
&self,
|
||||
input: RecordBatch,
|
||||
skip_fields: bool,
|
||||
skip_tags: bool,
|
||||
) -> Result<Option<RecordBatch>> {
|
||||
self.base.precise_filter_flat(input, skip_fields)
|
||||
self.base.precise_filter_flat(input, skip_fields, skip_tags)
|
||||
}
|
||||
|
||||
pub(crate) fn pre_filter_mode(&self) -> PreFilterMode {
|
||||
@@ -534,13 +548,16 @@ impl RangeBase {
|
||||
/// # Arguments
|
||||
/// * `input` - The RecordBatch to filter
|
||||
/// * `skip_fields` - Whether to skip field filters based on PreFilterMode
|
||||
/// * `skip_tags` - Whether to skip tag filters that were applied in an earlier phase
|
||||
pub(crate) fn precise_filter_flat(
|
||||
&self,
|
||||
input: RecordBatch,
|
||||
skip_fields: bool,
|
||||
skip_tags: bool,
|
||||
) -> Result<Option<RecordBatch>> {
|
||||
let mut tag_decode_state = TagDecodeState::new();
|
||||
let mask = self.compute_filter_mask_flat(&input, skip_fields, &mut tag_decode_state)?;
|
||||
let mask =
|
||||
self.compute_filter_mask_flat(&input, skip_fields, skip_tags, &mut tag_decode_state)?;
|
||||
|
||||
// If mask is None, the entire batch is filtered out
|
||||
let Some(mut mask) = mask else {
|
||||
@@ -558,9 +575,15 @@ impl RangeBase {
|
||||
mask = mask.bitand(&partition_mask);
|
||||
}
|
||||
|
||||
if mask.count_set_bits() == 0 {
|
||||
let num_selected = mask.count_set_bits();
|
||||
if num_selected == 0 {
|
||||
return Ok(None);
|
||||
}
|
||||
if num_selected == input.num_rows() {
|
||||
// Nothing was filtered out, e.g. all filters were skipped by
|
||||
// `skip_fields`/`skip_tags`. Avoid copying the whole batch.
|
||||
return Ok(Some(input));
|
||||
}
|
||||
|
||||
let filtered_batch =
|
||||
datatypes::arrow::compute::filter_record_batch(&input, &BooleanArray::from(mask))
|
||||
@@ -582,10 +605,12 @@ impl RangeBase {
|
||||
/// # Arguments
|
||||
/// * `input` - The RecordBatch to compute mask for
|
||||
/// * `skip_fields` - Whether to skip field filters based on PreFilterMode
|
||||
/// * `skip_tags` - Whether to skip tag filters that were applied in an earlier phase
|
||||
pub(crate) fn compute_filter_mask_flat(
|
||||
&self,
|
||||
input: &RecordBatch,
|
||||
skip_fields: bool,
|
||||
skip_tags: bool,
|
||||
tag_decode_state: &mut TagDecodeState,
|
||||
) -> Result<Option<BooleanBuffer>> {
|
||||
let mut mask = BooleanBuffer::new_set(input.num_rows());
|
||||
@@ -606,6 +631,9 @@ impl RangeBase {
|
||||
if skip_fields && filter_ctx.semantic_type() == SemanticType::Field {
|
||||
continue;
|
||||
}
|
||||
if skip_tags && filter_ctx.semantic_type() == SemanticType::Tag {
|
||||
continue;
|
||||
}
|
||||
|
||||
// Get the column directly by its projected index.
|
||||
// If the column is missing and it's not a tag/time column, this filter is skipped.
|
||||
@@ -783,7 +811,11 @@ mod tests {
|
||||
use std::sync::Arc;
|
||||
|
||||
use datafusion_expr::{col, lit};
|
||||
use datatypes::prelude::ConcreteDataType;
|
||||
use datatypes::schema::ColumnSchema;
|
||||
use datatypes::value::Value;
|
||||
use parquet::arrow::arrow_reader::RowSelector;
|
||||
use partition::expr::col as partition_col;
|
||||
|
||||
use super::*;
|
||||
use crate::read::read_columns::ReadColumns;
|
||||
@@ -829,10 +861,33 @@ mod tests {
|
||||
let batch = new_record_batch_with_custom_sequence(&["b", "x"], 0, 4, 1);
|
||||
|
||||
let mask = base
|
||||
.compute_filter_mask_flat(&batch, false, &mut TagDecodeState::new())
|
||||
.compute_filter_mask_flat(&batch, false, false, &mut TagDecodeState::new())
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(mask.count_set_bits(), 0);
|
||||
|
||||
let mask = base
|
||||
.compute_filter_mask_flat(&batch, false, true, &mut TagDecodeState::new())
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(mask.count_set_bits(), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_precise_filter_flat_returns_input_when_nothing_is_filtered() {
|
||||
let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
|
||||
let tag_filter =
|
||||
SimpleFilterContext::new_opt(&metadata, None, &col("tag_0").eq(lit("z"))).unwrap();
|
||||
let base = new_test_range_base(vec![tag_filter]);
|
||||
let batch = new_record_batch_with_custom_sequence(&["b", "x"], 0, 4, 1);
|
||||
|
||||
// The only filter is a tag filter, and it is skipped, so the batch must
|
||||
// come back untouched rather than being copied through `filter_record_batch`.
|
||||
let filtered = base
|
||||
.precise_filter_flat(batch.clone(), false, true)
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(batch, filtered);
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -859,12 +914,51 @@ mod tests {
|
||||
let batch = new_record_batch_with_custom_sequence(&["b", "x"], 0, 4, 1);
|
||||
|
||||
let mask = base
|
||||
.compute_filter_mask_flat(&batch, false, &mut TagDecodeState::new())
|
||||
.compute_filter_mask_flat(&batch, false, false, &mut TagDecodeState::new())
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(mask.count_set_bits(), 4);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_precise_filter_flat_applies_partition_filter_when_skipping_tags() {
|
||||
let metadata: RegionMetadataRef = Arc::new(sst_region_metadata());
|
||||
let tag_filter =
|
||||
SimpleFilterContext::new_opt(&metadata, None, &col("tag_0").eq(lit("z"))).unwrap();
|
||||
let mut base = new_test_range_base(vec![tag_filter]);
|
||||
let batch = new_record_batch_with_custom_sequence(&["b", "x"], 0, 4, 1);
|
||||
|
||||
let batch_schema = batch.schema();
|
||||
let tag_field = batch_schema.field(0);
|
||||
let partition_schema = Arc::new(Schema::new(vec![ColumnSchema::new(
|
||||
"tag_0".to_string(),
|
||||
ConcreteDataType::from_arrow_type(tag_field.data_type()),
|
||||
tag_field.is_nullable(),
|
||||
)]));
|
||||
let partition_expr = partition_col("tag_0")
|
||||
.gt_eq(Value::String("a".into()))
|
||||
.and(partition_col("tag_0").lt(Value::String("c".into())));
|
||||
base.partition_filter = Some(PartitionFilterContext {
|
||||
region_partition_physical_expr: partition_expr
|
||||
.try_as_physical_expr(partition_schema.arrow_schema())
|
||||
.unwrap(),
|
||||
partition_schema,
|
||||
});
|
||||
|
||||
let filtered = base
|
||||
.precise_filter_flat(batch, false, true)
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(filtered.num_rows(), 4);
|
||||
|
||||
let out_of_partition = new_record_batch_with_custom_sequence(&["z", "x"], 0, 4, 1);
|
||||
assert!(
|
||||
base.precise_filter_flat(out_of_partition, false, true)
|
||||
.unwrap()
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_refine_primary_key_selection_intersects_original_selection() {
|
||||
let original = Some(RowSelection::from(vec![
|
||||
|
||||
@@ -268,13 +268,14 @@ pub(crate) struct BulkFilterPlan {
|
||||
/// omitted. Callers use this as a pruning filter and must preserve the full predicate for
|
||||
/// authoritative filtering later in the scan.
|
||||
pub(crate) fn build_primary_key_filter(
|
||||
metadata: &RegionMetadataRef,
|
||||
sst_metadata: &RegionMetadataRef,
|
||||
expected_metadata: Option<&RegionMetadata>,
|
||||
predicate: Option<&Predicate>,
|
||||
) -> Option<CachedPrimaryKeyFilter> {
|
||||
let filters = predicate
|
||||
.into_iter()
|
||||
.flat_map(|predicate| predicate.exprs())
|
||||
.filter_map(|expr| SimpleFilterContext::new_opt(metadata, None, expr))
|
||||
.filter_map(|expr| SimpleFilterContext::new_opt(sst_metadata, expected_metadata, expr))
|
||||
.filter_map(|filter_ctx| {
|
||||
(filter_ctx.semantic_type() == SemanticType::Tag)
|
||||
.then(|| filter_ctx.filter().as_filter().cloned())
|
||||
@@ -285,8 +286,8 @@ pub(crate) fn build_primary_key_filter(
|
||||
return None;
|
||||
}
|
||||
|
||||
let codec = build_primary_key_codec(metadata.as_ref());
|
||||
let filter = codec.primary_key_filter(metadata, Arc::new(filters));
|
||||
let codec = build_primary_key_codec(sst_metadata.as_ref());
|
||||
let filter = codec.primary_key_filter(sst_metadata, Arc::new(filters));
|
||||
Some(CachedPrimaryKeyFilter::new(filter))
|
||||
}
|
||||
|
||||
@@ -370,13 +371,15 @@ pub(crate) fn build_bulk_filter_plan(
|
||||
/// reader always re-applies the original predicate, so the prefilter pass is purely a
|
||||
/// pruning hint.
|
||||
///
|
||||
/// Tag and timestamp predicates that lower to [`SimpleFilterEvaluator`] are an
|
||||
/// exception — the engine enforces them precisely, so the prefilter pass is the only
|
||||
/// place they execute. They are never silently dropped.
|
||||
/// With predicate prefiltering enabled, tag and timestamp predicates that lower to
|
||||
/// [`SimpleFilterEvaluator`] are an exception — the engine enforces them precisely in
|
||||
/// the prefilter pass. When it is disabled, all simple filters remain on the normal
|
||||
/// precise-filter path instead.
|
||||
pub(crate) fn build_reader_filter_plan(
|
||||
predicate: Option<&Predicate>,
|
||||
expected_metadata: Option<&RegionMetadata>,
|
||||
pre_filter_mode: PreFilterMode,
|
||||
enable_predicate_prefilter: bool,
|
||||
read_format: &FlatReadFormat,
|
||||
codec: &Arc<dyn PrimaryKeyCodec>,
|
||||
) -> ReaderFilterPlan {
|
||||
@@ -417,6 +420,11 @@ pub(crate) fn build_reader_filter_plan(
|
||||
// Prefer cheap simple filters first. They also preserve `Matched` /
|
||||
// `Pruned` states for columns that only exist in expected metadata.
|
||||
if let Some(filter_ctx) = SimpleFilterContext::new_opt(metadata, expected_metadata, expr) {
|
||||
if !enable_predicate_prefilter {
|
||||
remaining_simple_filters.push(filter_ctx);
|
||||
continue;
|
||||
}
|
||||
|
||||
// `Matched` and `Pruned` come from expected-metadata compatibility
|
||||
// and must stay in the main filter list so later phases keep that
|
||||
// outcome.
|
||||
@@ -452,6 +460,10 @@ pub(crate) fn build_reader_filter_plan(
|
||||
continue;
|
||||
}
|
||||
|
||||
if !enable_predicate_prefilter {
|
||||
continue;
|
||||
}
|
||||
|
||||
// Best-effort physical-filter prefilter (see fn-level doc): `new_opt`
|
||||
// returning `None` means the column is not in the projected arrow
|
||||
// schema, and dropping the predicate is safe because the upper
|
||||
@@ -464,6 +476,13 @@ pub(crate) fn build_reader_filter_plan(
|
||||
}
|
||||
}
|
||||
|
||||
if !enable_predicate_prefilter {
|
||||
return ReaderFilterPlan {
|
||||
remaining_simple_filters,
|
||||
prefilter_builder: None,
|
||||
};
|
||||
}
|
||||
|
||||
let pk_filter_expr_strs = (!pk_filter_contexts.is_empty()).then(|| {
|
||||
let mut expr_strs = pk_filter_contexts
|
||||
.iter()
|
||||
@@ -1153,6 +1172,7 @@ mod tests {
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
|
||||
use common_recordbatch::filter::SimpleFilterEvaluator;
|
||||
use datafusion_common::ScalarValue;
|
||||
use datafusion_expr::{col, lit};
|
||||
use datatypes::arrow::array::{
|
||||
ArrayRef, DictionaryArray, TimestampMillisecondArray, UInt8Array, UInt32Array, UInt64Array,
|
||||
@@ -1557,6 +1577,7 @@ mod tests {
|
||||
])),
|
||||
None,
|
||||
PreFilterMode::SkipFields,
|
||||
true,
|
||||
&full_read_format,
|
||||
&codec,
|
||||
);
|
||||
@@ -1584,6 +1605,7 @@ mod tests {
|
||||
Some(&Predicate::new(vec![col("tag_0").eq(lit("a"))])),
|
||||
None,
|
||||
PreFilterMode::All,
|
||||
true,
|
||||
&projected_read_format,
|
||||
&metric_codec,
|
||||
);
|
||||
@@ -1597,6 +1619,24 @@ mod tests {
|
||||
.is_some()
|
||||
);
|
||||
assert!(pk_prefilter_plan.remaining_simple_filters.is_empty());
|
||||
|
||||
let disabled_plan = build_reader_filter_plan(
|
||||
Some(&Predicate::new(vec![
|
||||
col("tag_0").eq(lit("a")),
|
||||
col("field_0").gt(lit(1_u64)),
|
||||
col("ts").gt_eq(lit(ScalarValue::TimestampMillisecond(Some(1), None))),
|
||||
])),
|
||||
None,
|
||||
PreFilterMode::All,
|
||||
false,
|
||||
&projected_read_format,
|
||||
&metric_codec,
|
||||
);
|
||||
assert!(disabled_plan.prefilter_builder.is_none());
|
||||
assert_eq!(
|
||||
remaining_simple_filter_columns(&disabled_plan.remaining_simple_filters),
|
||||
vec!["tag_0", "field_0", "ts"]
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -1622,6 +1662,7 @@ mod tests {
|
||||
Some(&Predicate::new(vec![expr_a.clone(), expr_b.clone()])),
|
||||
None,
|
||||
PreFilterMode::All,
|
||||
true,
|
||||
&read_format,
|
||||
&codec,
|
||||
);
|
||||
@@ -1629,6 +1670,7 @@ mod tests {
|
||||
Some(&Predicate::new(vec![expr_b, expr_a])),
|
||||
None,
|
||||
PreFilterMode::All,
|
||||
true,
|
||||
&read_format,
|
||||
&codec,
|
||||
);
|
||||
|
||||
@@ -85,7 +85,7 @@ use crate::sst::parquet::format::{INTERNAL_COLUMN_NUM, need_override_sequence};
|
||||
use crate::sst::parquet::json_align::{NestedSchemaAligner, ProjectedRecordBatchStream};
|
||||
use crate::sst::parquet::metadata::MetadataLoader;
|
||||
use crate::sst::parquet::prefilter::{
|
||||
CachedPrimaryKeyFilter, PrefilterContextBuilder, build_reader_filter_plan, execute_prefilter,
|
||||
PrefilterContextBuilder, build_reader_filter_plan, execute_prefilter,
|
||||
};
|
||||
use crate::sst::parquet::push_decoder::{
|
||||
SstParquetRangeFetcher, build_sst_parquet_record_batch_stream,
|
||||
@@ -183,6 +183,8 @@ pub struct ParquetReaderBuilder {
|
||||
compaction: bool,
|
||||
/// Mode to pre-filter columns.
|
||||
pre_filter_mode: PreFilterMode,
|
||||
/// Whether to run the reduced-column predicate prefilter pass.
|
||||
enable_predicate_prefilter: bool,
|
||||
/// Whether to decode primary key values eagerly when reading primary key format SSTs.
|
||||
decode_primary_key_values: bool,
|
||||
page_index_policy: PageIndexPolicy,
|
||||
@@ -217,6 +219,7 @@ impl ParquetReaderBuilder {
|
||||
expected_metadata: None,
|
||||
compaction: false,
|
||||
pre_filter_mode: PreFilterMode::All,
|
||||
enable_predicate_prefilter: true,
|
||||
decode_primary_key_values: false,
|
||||
page_index_policy: Default::default(),
|
||||
defer_optional_page_index: false,
|
||||
@@ -318,6 +321,13 @@ impl ParquetReaderBuilder {
|
||||
self
|
||||
}
|
||||
|
||||
/// Sets whether to run the reduced-column predicate prefilter pass.
|
||||
#[must_use]
|
||||
pub(crate) fn enable_predicate_prefilter(mut self, enable: bool) -> Self {
|
||||
self.enable_predicate_prefilter = enable;
|
||||
self
|
||||
}
|
||||
|
||||
/// Decodes primary key values eagerly when reading primary key format SSTs.
|
||||
#[must_use]
|
||||
pub(crate) fn decode_primary_key_values(mut self, decode: bool) -> Self {
|
||||
@@ -497,6 +507,7 @@ impl ParquetReaderBuilder {
|
||||
self.predicate.as_ref(),
|
||||
self.expected_metadata.as_deref(),
|
||||
self.pre_filter_mode,
|
||||
self.enable_predicate_prefilter,
|
||||
&read_format,
|
||||
&codec,
|
||||
);
|
||||
@@ -1838,13 +1849,6 @@ impl RowGroupReaderBuilder {
|
||||
self.prefilter_builder.is_some()
|
||||
}
|
||||
|
||||
/// Builds the encoded-primary-key filter selected by the reader filter plan.
|
||||
pub(crate) fn primary_key_filter(&self) -> Option<CachedPrimaryKeyFilter> {
|
||||
self.prefilter_builder
|
||||
.as_ref()
|
||||
.and_then(PrefilterContextBuilder::build_primary_key_filter)
|
||||
}
|
||||
|
||||
/// Builds a parquet record batch stream to read the row group at `row_group_idx`.
|
||||
///
|
||||
/// If prefiltering is applicable (based on `build_ctx`), this performs a two-phase read:
|
||||
@@ -1855,9 +1859,9 @@ impl RowGroupReaderBuilder {
|
||||
/// Predicates that cannot be lowered to prefilter columns (column not projected,
|
||||
/// expression not supported, etc.) are silently skipped. Correctness rests on the
|
||||
/// DataFusion `FilterExec` above this reader, which always re-applies the original
|
||||
/// predicate. Tag and timestamp predicates that flow through [`SimpleFilterEvaluator`]
|
||||
/// are an exception — the engine enforces them precisely, so the prefilter pass is the
|
||||
/// only place they execute. See [`build_reader_filter_plan`] for the bucketing rules.
|
||||
/// predicate. With predicate prefiltering enabled, tag and timestamp predicates that
|
||||
/// flow through [`SimpleFilterEvaluator`] are enforced precisely in this pass. See
|
||||
/// [`build_reader_filter_plan`] for the bucketing rules and disabled mode.
|
||||
///
|
||||
/// When the prefilter result selects no rows, the second read still issues but
|
||||
/// parquet-rs short-circuits before any column-chunk IO: the row-group state machine
|
||||
@@ -1907,7 +1911,6 @@ impl RowGroupReaderBuilder {
|
||||
///
|
||||
/// The series reader uses this after computing its own primary-key-only row
|
||||
/// selection.
|
||||
#[allow(dead_code)]
|
||||
pub(crate) async fn build_without_prefilter(
|
||||
&self,
|
||||
build_ctx: RowGroupBuildContext<'_>,
|
||||
|
||||
Reference in New Issue
Block a user