use std::any::Any;
use std::sync::Arc;
use async_trait::async_trait;
use datafusion::catalog::{Session, TableProvider};
use datafusion::common::Result as DfResult;
use datafusion::logical_expr::{
Expr, LogicalPlan, LogicalPlanBuilder, TableProviderFilterPushDown, TableType,
};
use datafusion::physical_plan::ExecutionPlan;
use crate::arrow::datatypes::SchemaRef;
use crate::tiering::store::ColdStore;
use super::provider::TieredTableProvider;
use super::version::{self, Resolution};
pub struct ResolvedTableProvider {
raw: Arc<TieredTableProvider>,
cold: Arc<dyn ColdStore>,
table: String,
resolution: LogicalPlan,
schema: SchemaRef,
}
impl std::fmt::Debug for ResolvedTableProvider {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ResolvedTableProvider")
.field("table", &self.table)
.finish_non_exhaustive()
}
}
impl ResolvedTableProvider {
pub fn new(
raw: Arc<TieredTableProvider>,
cold: Arc<dyn ColdStore>,
table: impl Into<String>,
resolution: LogicalPlan,
) -> Self {
let schema = Arc::new(resolution.schema().as_arrow().clone());
Self {
raw,
cold,
table: table.into(),
resolution,
schema,
}
}
fn resolution_without_stats(
&self,
watermark: crate::watermark::TieringWatermark,
filters: &[Expr],
) -> Option<Resolution> {
if self.raw.mode().forces_resolution() {
return Some(Resolution::Required);
}
let split = self.raw.split_at(watermark, filters);
if split.hot.is_some() {
return Some(Resolution::Required);
}
if split.cold.is_none() {
return Some(Resolution::Elided);
}
None
}
async fn resolution_from_stats(
&self,
watermark: crate::watermark::TieringWatermark,
filters: &[Expr],
) -> Resolution {
let Some(cold) = self.raw.split_at(watermark, filters).cold else {
return Resolution::Required;
};
let bounds = (
cold.start().unwrap_or(time::OffsetDateTime::UNIX_EPOCH),
cold.end().unwrap_or_else(|| {
time::OffsetDateTime::UNIX_EPOCH + time::Duration::days(365_000)
}),
);
match self.cold.version_stats(&self.table, bounds).await {
Ok(stats) => version::plan(&stats),
Err(_) => Resolution::Required,
}
}
}
#[async_trait]
impl TableProvider for ResolvedTableProvider {
fn as_any(&self) -> &dyn Any {
self
}
fn schema(&self) -> SchemaRef {
self.schema.clone()
}
fn table_type(&self) -> TableType {
TableType::View
}
async fn scan(
&self,
state: &dyn Session,
projection: Option<&Vec<usize>>,
filters: &[Expr],
limit: Option<usize>,
) -> DfResult<Arc<dyn ExecutionPlan>> {
let watermark = match super::PlannedWatermarks::of(state, &self.table) {
Some(pinned) => pinned,
None => self.raw.watermark().await?,
};
let (resolution, raw_scan) = match self.resolution_without_stats(watermark, filters) {
Some(decided) => (decided, None),
None => {
let scan = self
.raw
.scan_at(state, watermark, projection, filters, limit)
.await?;
(
self.resolution_from_stats(watermark, filters).await,
Some(scan),
)
}
};
let metrics = crate::observe::metrics();
let attrs = crate::observe::table(&self.table);
metrics.merge_elision_decisions.add(1, &attrs);
if resolution.is_elided() {
metrics.merge_elided.add(1, &attrs);
}
if resolution.is_elided() {
return match raw_scan {
Some(scan) => Ok(scan),
None => {
self.raw
.scan_at(state, watermark, projection, filters, limit)
.await
}
};
}
let mut builder = LogicalPlanBuilder::from(self.resolution.clone());
for filter in filters {
builder = builder.filter(filter.clone())?;
}
if let Some(indices) = projection {
let exprs = indices
.iter()
.map(|i| {
datafusion::logical_expr::col(datafusion::common::Column::from_name(
self.schema.field(*i).name(),
))
})
.collect::<Vec<_>>();
builder = builder.project(exprs)?;
}
if let Some(fetch) = limit {
builder = builder.limit(0, Some(fetch))?;
}
state.create_physical_plan(&builder.build()?).await
}
fn supports_filters_pushdown(
&self,
filters: &[&Expr],
) -> DfResult<Vec<TableProviderFilterPushDown>> {
Ok(vec![TableProviderFilterPushDown::Inexact; filters.len()])
}
}