1use std::cmp::Ordering;
16use std::ops::Range;
17
18use datafusion_expr::{col, lit};
19use datatypes::arrow::array::{Int64Array, UInt32Array, UInt64Array};
20use datatypes::arrow::datatypes::{DataType, SchemaRef};
21use futures::TryStreamExt;
22use object_store::ObjectStore;
23use snafu::{OptionExt, ensure};
24use table::predicate::Predicate;
25
26use crate::error::{InvalidRecordBatchSnafu, Result, UnexpectedSnafu};
27use crate::series_index::MetricSeriesId;
28use crate::sst::parquet::index_reader::ParquetIndexReader;
29use crate::sst::range_index::{
30 END_COLUMN, ROW_GROUP_ID_COLUMN, START_COLUMN, TABLE_ID_COLUMN, TSID_COLUMN,
31};
32
33pub struct SstRangeIndexSearcher {
35 reader: ParquetIndexReader,
36}
37
38impl SstRangeIndexSearcher {
39 pub async fn open(object_store: ObjectStore, path: &str) -> Result<Self> {
41 let reader = ParquetIndexReader::open(object_store, path).await?;
42 validate_index_schema(reader.schema())?;
43 Ok(Self { reader })
44 }
45
46 pub async fn search(
55 &self,
56 row_group_id: u32,
57 series: &[MetricSeriesId],
58 ) -> Result<Vec<Range<usize>>> {
59 if series.is_empty() {
60 return Ok(Vec::new());
61 }
62
63 validate_sorted_series(series)?;
64 let predicate = search_predicate(row_group_id, series)?;
65 let mut batches = self.reader.read(
66 &predicate,
67 &[
68 ROW_GROUP_ID_COLUMN,
69 TABLE_ID_COLUMN,
70 TSID_COLUMN,
71 START_COLUMN,
72 END_COLUMN,
73 ],
74 )?;
75 let mut merge = RangeMergeState::new(row_group_id, series);
76
77 while let Some(batch) = batches.try_next().await? {
78 if merge.append_batch(&batch)? {
79 break;
80 }
81 }
82
83 Ok(merge.finish())
84 }
85}
86
87fn validate_sorted_series(series: &[MetricSeriesId]) -> Result<()> {
88 if let Some(pair) = series.windows(2).find(|pair| pair[0] > pair[1]) {
89 return InvalidRecordBatchSnafu {
90 reason: format!(
91 "range index search series are not sorted: {:?} appears before {:?}",
92 pair[0], pair[1]
93 ),
94 }
95 .fail();
96 }
97 Ok(())
98}
99
100fn search_predicate(row_group_id: u32, series: &[MetricSeriesId]) -> Result<Predicate> {
101 let min_table_id = series
102 .first()
103 .context(UnexpectedSnafu {
104 reason: "cannot build a range-index predicate for an empty series set",
105 })?
106 .table_id;
107 let max_table_id = series
108 .last()
109 .context(UnexpectedSnafu {
110 reason: "cannot build a range-index predicate for an empty series set",
111 })?
112 .table_id;
113
114 Ok(Predicate::new(vec![
115 col(ROW_GROUP_ID_COLUMN).eq(lit(row_group_id)),
116 col(TABLE_ID_COLUMN).gt_eq(lit(min_table_id)),
117 col(TABLE_ID_COLUMN).lt_eq(lit(max_table_id)),
118 ]))
119}
120
121fn validate_index_schema(schema: &SchemaRef) -> Result<()> {
122 for (name, data_type) in [
123 (ROW_GROUP_ID_COLUMN, DataType::UInt32),
124 (TABLE_ID_COLUMN, DataType::UInt32),
125 (TSID_COLUMN, DataType::UInt64),
126 (START_COLUMN, DataType::Int64),
127 (END_COLUMN, DataType::Int64),
128 ] {
129 let field = schema
130 .field_with_name(name)
131 .ok()
132 .with_context(|| InvalidRecordBatchSnafu {
133 reason: format!("range index is missing column {name}"),
134 })?;
135 ensure!(
136 field.data_type() == &data_type && !field.is_nullable(),
137 InvalidRecordBatchSnafu {
138 reason: format!(
139 "range index column {name} must be non-nullable {data_type:?}, got {:?}",
140 field.data_type()
141 ),
142 }
143 );
144 }
145 Ok(())
146}
147
148struct RangeMergeState<'a> {
149 row_group_id: u32,
151 series: &'a [MetricSeriesId],
153 series_index: usize,
155 last_index_key: Option<(u32, MetricSeriesId)>,
157 ranges: Vec<Range<usize>>,
159}
160
161impl<'a> RangeMergeState<'a> {
162 fn new(row_group_id: u32, series: &'a [MetricSeriesId]) -> Self {
163 Self {
164 row_group_id,
165 series,
166 series_index: 0,
167 last_index_key: None,
168 ranges: Vec::new(),
169 }
170 }
171
172 fn append_batch(
174 &mut self,
175 batch: &datatypes::arrow::record_batch::RecordBatch,
176 ) -> Result<bool> {
177 let row_group_ids = typed_column::<UInt32Array>(batch, ROW_GROUP_ID_COLUMN, "UInt32")?;
178 let table_ids = typed_column::<UInt32Array>(batch, TABLE_ID_COLUMN, "UInt32")?;
179 let tsids = typed_column::<UInt64Array>(batch, TSID_COLUMN, "UInt64")?;
180 let starts = typed_column::<Int64Array>(batch, START_COLUMN, "Int64")?;
181 let ends = typed_column::<Int64Array>(batch, END_COLUMN, "Int64")?;
182
183 for row in 0..batch.num_rows() {
184 let index_series = MetricSeriesId {
185 table_id: table_ids.value(row),
186 tsid: tsids.value(row),
187 };
188 let index_key = (row_group_ids.value(row), index_series);
189 ensure!(
190 self.last_index_key.is_none_or(|last| last < index_key),
191 InvalidRecordBatchSnafu {
192 reason: format!(
193 "range index rows are not strictly sorted: {index_key:?} follows {:?}",
194 self.last_index_key
195 ),
196 }
197 );
198 self.last_index_key = Some(index_key);
199
200 match index_key.0.cmp(&self.row_group_id) {
201 Ordering::Less => continue,
202 Ordering::Greater => return Ok(true),
203 Ordering::Equal => {}
204 }
205
206 while self.series_index < self.series.len()
207 && self.series[self.series_index] < index_series
208 {
209 self.advance_series();
210 }
211 if self.series_index == self.series.len() {
212 return Ok(true);
213 }
214
215 match self.series[self.series_index].cmp(&index_series) {
216 Ordering::Less => {
217 return UnexpectedSnafu {
218 reason: "range-index merge cursor did not advance past a smaller series",
219 }
220 .fail();
221 }
222 Ordering::Greater => continue,
223 Ordering::Equal => {
224 self.append_range(starts.value(row), ends.value(row), row)?;
225 self.advance_series();
226 if self.series_index == self.series.len() {
227 return Ok(true);
228 }
229 }
230 }
231 }
232 Ok(false)
233 }
234
235 fn advance_series(&mut self) {
236 let current = self.series[self.series_index];
237 while self.series_index < self.series.len() && self.series[self.series_index] == current {
238 self.series_index += 1;
239 }
240 }
241
242 fn append_range(&mut self, start: i64, end: i64, row: usize) -> Result<()> {
243 let start = usize::try_from(start).map_err(|_| {
244 InvalidRecordBatchSnafu {
245 reason: format!("range index contains negative start offset at row {row}"),
246 }
247 .build()
248 })?;
249 let end = usize::try_from(end).map_err(|_| {
250 InvalidRecordBatchSnafu {
251 reason: format!("range index contains negative end offset at row {row}"),
252 }
253 .build()
254 })?;
255 ensure!(
256 start < end,
257 InvalidRecordBatchSnafu {
258 reason: format!("range index contains invalid range {start}..{end} at row {row}"),
259 }
260 );
261
262 if let Some(last) = self.ranges.last_mut() {
263 ensure!(
264 start >= last.end,
265 InvalidRecordBatchSnafu {
266 reason: format!(
267 "range index contains overlapping or unsorted range {start}..{end} after {}..{}",
268 last.start, last.end
269 ),
270 }
271 );
272 if start == last.end {
273 last.end = end;
274 return Ok(());
275 }
276 }
277 self.ranges.push(start..end);
278 Ok(())
279 }
280
281 fn finish(self) -> Vec<Range<usize>> {
282 self.ranges
283 }
284}
285
286fn typed_column<'a, T: 'static>(
287 batch: &'a datatypes::arrow::record_batch::RecordBatch,
288 name: &str,
289 data_type: &str,
290) -> Result<&'a T> {
291 let index = batch
292 .schema()
293 .index_of(name)
294 .ok()
295 .with_context(|| InvalidRecordBatchSnafu {
296 reason: format!("range index batch is missing column {name}"),
297 })?;
298 batch
299 .column(index)
300 .as_any()
301 .downcast_ref::<T>()
302 .with_context(|| InvalidRecordBatchSnafu {
303 reason: format!("range index column {name} is not {data_type}"),
304 })
305}
306
307#[cfg(test)]
308mod tests {
309 use std::sync::Arc;
310
311 use datatypes::arrow::array::{ArrayRef, BinaryArray};
312 use datatypes::arrow::datatypes::{Field, Schema};
313 use datatypes::arrow::record_batch::RecordBatch;
314 use object_store::services::Memory;
315 use store_api::codec::PrimaryKeyEncoding;
316 use store_api::metadata::RegionMetadataRef;
317 use store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME;
318
319 use super::*;
320 use crate::sst::range_index::{
321 SstRangeIndexWriter, SstRangeIndexWriterOptions, range_index_schema,
322 };
323 use crate::test_util::sst_util::{new_sparse_primary_key, sst_region_metadata_with_encoding};
324
325 fn object_store() -> ObjectStore {
326 ObjectStore::new(Memory::default()).unwrap()
327 }
328
329 fn series(table_id: u32, tsid: u64) -> MetricSeriesId {
330 MetricSeriesId { table_id, tsid }
331 }
332
333 fn primary_key_batch(metadata: &RegionMetadataRef, ids: &[(u32, u64)]) -> RecordBatch {
334 let primary_keys = ids
335 .iter()
336 .map(|(table_id, tsid)| new_sparse_primary_key(&["a", "x"], metadata, *table_id, *tsid))
337 .collect::<Vec<_>>();
338 let schema = Arc::new(Schema::new(vec![Field::new(
339 PRIMARY_KEY_COLUMN_NAME,
340 DataType::Binary,
341 false,
342 )]));
343 RecordBatch::try_new(
344 schema,
345 vec![Arc::new(BinaryArray::from_iter_values(
346 primary_keys.iter().map(Vec::as_slice),
347 ))],
348 )
349 .unwrap()
350 }
351
352 async fn write_index(store: &ObjectStore, path: &str) {
353 let metadata = Arc::new(sst_region_metadata_with_encoding(
354 PrimaryKeyEncoding::Sparse,
355 ));
356 let mut writer = SstRangeIndexWriter::try_new(
357 metadata.clone(),
358 store.clone(),
359 path,
360 SstRangeIndexWriterOptions {
361 index_row_group_size: 2,
362 },
363 )
364 .await
365 .unwrap();
366 writer
367 .write(
368 0,
369 &primary_key_batch(
370 &metadata,
371 &[(1, 10), (1, 10), (1, 20), (2, 10), (2, 20), (2, 20)],
372 ),
373 )
374 .await
375 .unwrap();
376 writer
377 .write(1, &primary_key_batch(&metadata, &[(2, 20), (2, 20)]))
378 .await
379 .unwrap();
380 writer.finish().await.unwrap();
381 }
382
383 #[tokio::test]
384 async fn search_filters_exact_series_pairs_and_coalesces_ranges() {
385 let store = object_store();
386 let path = "range-search.parquet";
387 write_index(&store, path).await;
388 let searcher = SstRangeIndexSearcher::open(store, path).await.unwrap();
389
390 let ranges = searcher
391 .search(0, &[series(1, 10), series(2, 20)])
392 .await
393 .unwrap();
394 assert_eq!(ranges, vec![0..2, 4..6]);
395
396 let ranges = searcher
397 .search(0, &[series(1, 10), series(1, 20)])
398 .await
399 .unwrap();
400 assert_eq!(ranges, vec![0..3]);
401
402 let ranges = searcher
403 .search(0, &[series(1, 15), series(2, 20)])
404 .await
405 .unwrap();
406 assert_eq!(ranges, vec![4..6]);
407
408 let ranges = searcher
409 .search(0, &[series(1, 10), series(2, 30)])
410 .await
411 .unwrap();
412 assert_eq!(ranges, vec![0..2]);
413
414 let ranges = searcher
415 .search(1, &[series(2, 20), series(2, 20)])
416 .await
417 .unwrap();
418 assert_eq!(ranges, vec![0..2]);
419
420 assert!(
421 searcher
422 .search(1, &[series(1, 10)])
423 .await
424 .unwrap()
425 .is_empty()
426 );
427
428 assert!(searcher.search(0, &[]).await.unwrap().is_empty());
429
430 let error = searcher
431 .search(0, &[series(2, 20), series(1, 10)])
432 .await
433 .unwrap_err();
434 assert!(error.to_string().contains("not sorted"), "{error}");
435 }
436
437 #[tokio::test]
438 async fn opening_a_missing_index_fails() {
439 assert!(
440 SstRangeIndexSearcher::open(object_store(), "does-not-exist.parquet")
441 .await
442 .is_err()
443 );
444 }
445
446 #[tokio::test]
447 async fn pruning_uses_the_source_row_group_and_table_id_range() {
448 let store = object_store();
449 let path = "range-pruning.parquet";
450 write_index(&store, path).await;
451 let reader = ParquetIndexReader::open(store, path).await.unwrap();
452 let predicate = search_predicate(0, &[series(1, 999), series(2, 999)]).unwrap();
453
454 assert_eq!(reader.row_groups_to_read(&predicate), vec![0, 1]);
455 }
456
457 #[test]
458 fn validates_schema_and_range_offsets() {
459 let nullable_schema = Arc::new(Schema::new(vec![
460 Field::new(ROW_GROUP_ID_COLUMN, DataType::UInt32, false),
461 Field::new(TABLE_ID_COLUMN, DataType::UInt32, false),
462 Field::new(TSID_COLUMN, DataType::UInt64, false),
463 Field::new(START_COLUMN, DataType::Int64, true),
464 Field::new(END_COLUMN, DataType::Int64, false),
465 ]));
466 assert!(validate_index_schema(&nullable_schema).is_err());
467
468 let batch = RecordBatch::try_new(
469 range_index_schema(),
470 vec![
471 Arc::new(UInt32Array::from(vec![0])) as ArrayRef,
472 Arc::new(UInt32Array::from(vec![1])),
473 Arc::new(UInt64Array::from(vec![10])),
474 Arc::new(Int64Array::from(vec![-1])),
475 Arc::new(Int64Array::from(vec![2])),
476 ],
477 )
478 .unwrap();
479 let selected = [series(1, 10)];
480 let mut merge = RangeMergeState::new(0, &selected);
481 assert!(merge.append_batch(&batch).is_err());
482
483 let unsorted_batch = RecordBatch::try_new(
484 range_index_schema(),
485 vec![
486 Arc::new(UInt32Array::from(vec![0, 0])) as ArrayRef,
487 Arc::new(UInt32Array::from(vec![1, 1])),
488 Arc::new(UInt64Array::from(vec![20, 10])),
489 Arc::new(Int64Array::from(vec![0, 1])),
490 Arc::new(Int64Array::from(vec![1, 2])),
491 ],
492 )
493 .unwrap();
494 let selected = [series(1, 20), series(1, 30)];
495 let mut merge = RangeMergeState::new(0, &selected);
496 assert!(merge.append_batch(&unsorted_batch).is_err());
497
498 let make_batch = |tsid, start, end| {
499 RecordBatch::try_new(
500 range_index_schema(),
501 vec![
502 Arc::new(UInt32Array::from(vec![0])) as ArrayRef,
503 Arc::new(UInt32Array::from(vec![1])),
504 Arc::new(UInt64Array::from(vec![tsid])),
505 Arc::new(Int64Array::from(vec![start])),
506 Arc::new(Int64Array::from(vec![end])),
507 ],
508 )
509 .unwrap()
510 };
511 let selected = [series(1, 10), series(1, 20)];
512 let mut merge = RangeMergeState::new(0, &selected);
513 assert!(!merge.append_batch(&make_batch(10, 0, 1)).unwrap());
514 assert!(merge.append_batch(&make_batch(20, 1, 2)).unwrap());
515 assert_eq!(merge.finish(), vec![0..2]);
516 }
517}