use std::sync::Arc;
use arrow::array::{RecordBatch, RecordBatchOptions};
use arrow::datatypes::SchemaRef;
use datafusion_common::Result;
use datafusion_physical_expr::projection::{ProjectionExprs, Projector};
use datafusion_physical_expr::utils::reassign_expr_columns;
use datafusion_physical_expr_adapter::replace_columns_with_literals;
use parquet::arrow::ProjectionMask;
use parquet::schema::types::SchemaDescriptor;
use crate::opener::{VirtualColumnsState, append_fields};
use crate::projection_read_plan::build_projection_read_plan;
pub(crate) struct DecoderProjection {
projection_mask: ProjectionMask,
projector: Projector,
output_schema: SchemaRef,
replace_schema: bool,
}
impl DecoderProjection {
pub(crate) fn try_new(
projection: &ProjectionExprs,
physical_file_schema: &SchemaRef,
parquet_schema: &SchemaDescriptor,
output_schema: &SchemaRef,
virtual_state: Option<&VirtualColumnsState>,
) -> Result<Self> {
let projection_for_read_plan = match virtual_state {
None => projection.clone(),
Some(state) => projection.clone().try_map_exprs(|expr| {
replace_columns_with_literals(expr, state.null_replacements())
})?,
};
let read_plan = build_projection_read_plan(
projection_for_read_plan.expr_iter(),
physical_file_schema,
parquet_schema,
);
let stream_schema = match virtual_state {
Some(state) => {
append_fields(&read_plan.projected_schema, state.virtual_columns())
}
None => Arc::clone(&read_plan.projected_schema),
};
let rebased_projection = projection
.clone()
.try_map_exprs(|expr| reassign_expr_columns(expr, &stream_schema))?;
let projector = rebased_projection.make_projector(&stream_schema)?;
let replace_schema = projector.output_schema() != output_schema;
Ok(Self {
projection_mask: read_plan.projection_mask,
projector,
output_schema: Arc::clone(output_schema),
replace_schema,
})
}
pub(crate) fn projection_mask(&self) -> &ProjectionMask {
&self.projection_mask
}
pub(crate) fn map(&self, batch: &RecordBatch) -> Result<RecordBatch> {
let projected = self.projector.project_batch(batch)?;
if !self.replace_schema {
return Ok(projected);
}
let (_stream_schema, arrays, num_rows) = projected.into_parts();
let options = RecordBatchOptions::new().with_row_count(Some(num_rows));
Ok(RecordBatch::try_new_with_options(
Arc::clone(&self.output_schema),
arrays,
&options,
)?)
}
}