Skip to content

Commit 40112f4

Browse files
authored
Actually wire the pluggable expression convertor (#7730)
## Summary - Closes #7731 Turns out I didn't wire the expression convertor extension point correctly. This PR both fixes that AND adds a test to make sure this behavior is maintained. Signed-off-by: Adam Gutglick <adam@spiraldb.com>
1 parent c73dbb2 commit 40112f4

1 file changed

Lines changed: 99 additions & 6 deletions

File tree

vortex-datafusion/src/persistent/source.rs

Lines changed: 99 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -301,15 +301,13 @@ impl VortexSource {
301301
self.options = opts;
302302
self
303303
}
304-
}
305304

306-
impl FileSource for VortexSource {
307-
fn create_file_opener(
305+
fn create_vortex_opener(
308306
&self,
309307
object_store: Arc<dyn ObjectStore>,
310308
base_config: &FileScanConfig,
311309
partition: usize,
312-
) -> DFResult<Arc<dyn FileOpener>> {
310+
) -> DFResult<VortexOpener> {
313311
let batch_size = self
314312
.batch_size
315313
.vortex_expect("batch_size must be supplied to VortexSource");
@@ -339,13 +337,28 @@ impl FileSource for VortexSource {
339337
layout_readers: Arc::clone(&self.layout_readers),
340338
natural_split_ranges: Arc::clone(&self.natural_split_ranges),
341339
has_output_ordering: !base_config.output_ordering.is_empty(),
342-
expression_convertor: Arc::new(DefaultExpressionConvertor::default()),
340+
expression_convertor: Arc::clone(&self.expression_convertor),
343341
file_metadata_cache: self.file_metadata_cache.clone(),
344342
projection_pushdown: self.options.projection_pushdown,
345343
scan_concurrency: self.options.scan_concurrency,
346344
};
347345

348-
Ok(Arc::new(opener))
346+
Ok(opener)
347+
}
348+
}
349+
350+
impl FileSource for VortexSource {
351+
fn create_file_opener(
352+
&self,
353+
object_store: Arc<dyn ObjectStore>,
354+
base_config: &FileScanConfig,
355+
partition: usize,
356+
) -> DFResult<Arc<dyn FileOpener>> {
357+
Ok(Arc::new(self.create_vortex_opener(
358+
object_store,
359+
base_config,
360+
partition,
361+
)?))
349362
}
350363

351364
fn as_any(&self) -> &dyn Any {
@@ -477,3 +490,83 @@ impl FileSource for VortexSource {
477490
&self.table_schema
478491
}
479492
}
493+
494+
#[cfg(test)]
495+
mod tests {
496+
use arrow_schema::DataType;
497+
use arrow_schema::Field;
498+
use arrow_schema::Schema;
499+
use datafusion_datasource::file_scan_config::FileScanConfigBuilder;
500+
use datafusion_execution::object_store::ObjectStoreUrl;
501+
use object_store::memory::InMemory;
502+
use vortex::VortexSessionDefault;
503+
504+
use super::*;
505+
use crate::convert::exprs::ProcessedProjection;
506+
507+
struct TrackingExpressionConvertor {
508+
inner: DefaultExpressionConvertor,
509+
}
510+
511+
impl ExpressionConvertor for TrackingExpressionConvertor {
512+
fn can_be_pushed_down(&self, expr: &PhysicalExprRef, schema: &Schema) -> bool {
513+
self.inner.can_be_pushed_down(expr, schema)
514+
}
515+
516+
fn convert(&self, expr: &dyn PhysicalExpr) -> DFResult<vortex::expr::Expression> {
517+
self.inner.convert(expr)
518+
}
519+
520+
fn split_projection(
521+
&self,
522+
source_projection: ProjectionExprs,
523+
input_schema: &Schema,
524+
output_schema: &Schema,
525+
) -> DFResult<ProcessedProjection> {
526+
self.inner
527+
.split_projection(source_projection, input_schema, output_schema)
528+
}
529+
530+
fn no_pushdown_projection(
531+
&self,
532+
source_projection: ProjectionExprs,
533+
input_schema: &Schema,
534+
) -> DFResult<ProcessedProjection> {
535+
self.inner
536+
.no_pushdown_projection(source_projection, input_schema)
537+
}
538+
}
539+
540+
#[test]
541+
fn create_vortex_opener_preserves_expression_convertor() -> anyhow::Result<()> {
542+
let file_schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)]));
543+
let expression_convertor = Arc::new(TrackingExpressionConvertor {
544+
inner: DefaultExpressionConvertor::default(),
545+
}) as Arc<dyn ExpressionConvertor>;
546+
547+
let mut source = VortexSource::new(
548+
TableSchema::from_file_schema(file_schema),
549+
VortexSession::default(),
550+
)
551+
.with_expression_convertor(Arc::clone(&expression_convertor));
552+
source.batch_size = Some(100);
553+
554+
let config = FileScanConfigBuilder::new(
555+
ObjectStoreUrl::local_filesystem(),
556+
Arc::new(source.clone()),
557+
)
558+
.build();
559+
560+
let opener = source.create_vortex_opener(
561+
Arc::new(InMemory::new()) as Arc<dyn ObjectStore>,
562+
&config,
563+
0,
564+
)?;
565+
566+
assert!(Arc::ptr_eq(
567+
&opener.expression_convertor,
568+
&expression_convertor
569+
));
570+
Ok(())
571+
}
572+
}

0 commit comments

Comments
 (0)