use std::fs::File;
use std::sync::Arc;
use crate::error::{Error, Result};
use crate::schema::{LogicalType, Schema};
use super::block::BlockMetadata;
use super::budget::{ScanBudget, ScanMetrics};
use super::reader::Reader;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PrimaryRange {
Timestamp {
start: i64,
end: i64,
},
Date32 {
start: i32,
end: i32,
},
}
impl PrimaryRange {
pub fn timestamp(start: i64, end: i64) -> Self {
Self::Timestamp { start, end }
}
pub fn date32(start: i32, end: i32) -> Self {
Self::Date32 { start, end }
}
fn bounds(self) -> (i64, i64) {
match self {
Self::Timestamp { start, end } => (start, end),
Self::Date32 { start, end } => (i64::from(start), i64::from(end)),
}
}
}
#[derive(Debug)]
pub(crate) struct ScanPlan {
pub(crate) projection: Vec<usize>,
pub(crate) output_schema: Arc<Schema>,
pub(crate) range: Option<PrimaryRange>,
}
impl ScanPlan {
pub(crate) fn new(schema: &Schema) -> Self {
let projection: Vec<usize> = (0..schema.column_count()).collect();
let output_schema = projected_schema(schema, &projection);
Self {
projection,
output_schema,
range: None,
}
}
pub(crate) fn project<I, S>(&mut self, schema: &Schema, columns: I) -> Result<()>
where
I: IntoIterator<Item = S>,
S: AsRef<str>,
{
let mut projection = Vec::new();
for name in columns {
let name = name.as_ref();
let index = schema
.columns()
.iter()
.position(|column| column.name() == name)
.ok_or_else(|| {
Error::invalid_argument(format!("unknown projected column {name}"))
})?;
if projection.contains(&index) {
return Err(Error::invalid_argument(format!(
"projected column {name} was requested more than once"
)));
}
projection.push(index);
}
self.output_schema = projected_schema(schema, &projection);
self.projection = projection;
Ok(())
}
pub(crate) fn primary_range(&mut self, schema: &Schema, range: PrimaryRange) -> Result<()> {
let (start, end) = range.bounds();
if start > end {
return Err(Error::invalid_argument(
"primary range start must not be greater than its end",
));
}
let primary = schema.primary_column().ok_or_else(|| {
Error::invalid_argument("primary ranges require a schema primary column")
})?;
let compatible = matches!(
(range, primary.logical_type()),
(
PrimaryRange::Timestamp { .. },
LogicalType::Timestamp { .. }
) | (PrimaryRange::Date32 { .. }, LogicalType::Date32)
);
if !compatible {
return Err(Error::invalid_argument(format!(
"primary range type does not match primary column {}",
primary.name()
)));
}
self.range = Some(range);
Ok(())
}
fn internal_primary(&self, schema: &Schema) -> Option<usize> {
super::decode::primary_index(schema)
.filter(|index| self.range.is_some() || self.projection.contains(index))
}
pub(crate) fn selection<'a>(
&'a self,
source_schema: &'a Arc<Schema>,
) -> super::decode::Selection<'a> {
super::decode::Selection {
source_schema: Arc::clone(source_schema),
output_schema: Arc::clone(&self.output_schema),
selected: &self.projection,
primary: self.internal_primary(source_schema),
}
}
pub(crate) fn should_prune(&self, block: &BlockMetadata) -> bool {
let Some(range) = self.range else {
return false;
};
let (start, end) = range.bounds();
if start == end {
return true;
}
let Some(bounds) = block.primary_bounds() else {
return false;
};
bounds.max() < start || bounds.min() >= end
}
pub(crate) fn filter_batch(
&self,
primary_sorted: bool,
batch: crate::RecordBatch,
primary_values: &[i64],
range: PrimaryRange,
) -> Result<Option<crate::RecordBatch>> {
let (start, end) = range.bounds();
if start == end {
return Ok(None);
}
if primary_sorted {
let first = primary_values.partition_point(|value| *value < start);
let last = primary_values.partition_point(|value| *value < end);
if first == last {
return Ok(None);
}
return Ok(Some(batch.slice(first, last)));
}
let mut indices = Vec::new();
indices.try_reserve(primary_values.len()).map_err(|_| {
Error::resource_limit("unable to allocate unsorted range selection", None)
})?;
indices.extend(
primary_values
.iter()
.enumerate()
.filter_map(|(index, value)| (*value >= start && *value < end).then_some(index)),
);
if indices.is_empty() {
return Ok(None);
}
Ok(Some(batch.take(&indices)))
}
}
#[derive(Debug)]
pub struct Scan<'reader> {
pub(crate) reader: &'reader Reader,
pub(crate) next_index: usize,
pub(crate) file: Option<File>,
pub(crate) plan: ScanPlan,
pub(crate) budget: ScanBudget,
}
impl<'reader> Scan<'reader> {
pub(crate) fn new(reader: &'reader Reader) -> Self {
Self {
reader,
next_index: 0,
file: None,
plan: ScanPlan::new(reader.schema()),
budget: ScanBudget::new(
reader.limits().max_rows_per_scan(),
reader.limits().max_decoded_scan_bytes(),
),
}
}
pub fn project<I, S>(mut self, columns: I) -> Result<Self>
where
I: IntoIterator<Item = S>,
S: AsRef<str>,
{
self.plan.project(self.reader.schema(), columns)?;
Ok(self)
}
pub fn primary_range(mut self, range: PrimaryRange) -> Result<Self> {
self.plan.primary_range(self.reader.schema(), range)?;
Ok(self)
}
pub fn file_order(self) -> Self {
self
}
pub fn remaining_candidate_blocks(&self) -> usize {
self.remaining_candidates()
}
pub fn metrics(&self) -> ScanMetrics {
self.budget.metrics()
}
fn remaining_candidates(&self) -> usize {
self.reader.blocks()[self.next_index..]
.iter()
.filter(|block| !self.plan.should_prune(block))
.count()
}
}
impl Iterator for Scan<'_> {
type Item = Result<crate::RecordBatch>;
fn next(&mut self) -> Option<Self::Item> {
while self.next_index < self.reader.blocks().len() {
let index = self.next_index;
self.next_index += 1;
let block = &self.reader.blocks()[index];
let pruned = self.plan.should_prune(block);
if let Err(error) = self.budget.record_block(pruned) {
return Some(Err(error));
}
if pruned {
continue;
}
let selection = self.plan.selection(self.reader.schema_handle());
let file = match &mut self.file {
Some(file) => file,
None => match File::open(self.reader.path()) {
Ok(file) => self.file.insert(file),
Err(error) => {
return Some(Err(
Error::io(error, None).with_context(crate::ErrorContext::File)
));
}
},
};
if let Err(error) = self.budget.charge_rows(block.row_count()) {
return Some(Err(error));
}
let decoded =
self.reader
.decode_selected_block_at(file, index, &selection, &mut self.budget);
match decoded {
Err(error) => return Some(Err(error)),
Ok(decoded) => {
let super::decode::DecodedBlock {
batch,
primary_values,
primary_sorted,
} = decoded;
let Some(range) = self.plan.range else {
if let Err(error) = self.budget.record_rows(batch.row_count()) {
return Some(Err(error));
}
return Some(Ok(batch));
};
let Some(primary_values) = primary_values else {
return Some(Err(Error::internal(
"range scan did not decode its primary column",
)));
};
match self
.plan
.filter_batch(primary_sorted, batch, &primary_values, range)
{
Ok(Some(batch)) => {
if let Err(error) = self.budget.record_rows(batch.row_count()) {
return Some(Err(error));
}
return Some(Ok(batch));
}
Ok(None) => continue,
Err(error) => return Some(Err(error)),
}
}
}
}
None
}
fn size_hint(&self) -> (usize, Option<usize>) {
(0, Some(self.remaining_candidates()))
}
}
impl std::iter::FusedIterator for Scan<'_> {}
fn projected_schema(schema: &Schema, projection: &[usize]) -> Arc<Schema> {
let columns = projection
.iter()
.map(|&index| schema.columns()[index].clone())
.collect();
let primary = schema.primary_column_id().filter(|id| {
projection
.iter()
.any(|&index| schema.columns()[index].id() == *id)
});
Arc::new(Schema::new(schema.schema_id(), columns, primary))
}