Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 5 additions & 3 deletions native/core/src/execution/planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,7 @@ use datafusion_comet_proto::{
},
spark_partitioning::{partitioning::PartitioningStruct, Partitioning as SparkPartitioning},
};
use datafusion_comet_spark_expr::parquet_support::SparkParquetOptions;
use datafusion_comet_spark_expr::{
ArrayInsert, Avg, AvgDecimal, BitwiseNotExpr, Cast, CheckOverflow, Contains, Correlation,
Covariance, CreateNamedStruct, DateTruncExpr, EndsWith, GetArrayStructFields, GetStructField,
Expand Down Expand Up @@ -1156,13 +1157,14 @@ impl PhysicalPlanner {
table_parquet_options.global.pushdown_filters = true;
table_parquet_options.global.reorder_filters = true;

let mut spark_cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC", false);
spark_cast_options.allow_cast_unsigned_ints = true;
let mut spark_parquet_options =
SparkParquetOptions::new(EvalMode::Legacy, "UTC", false);
spark_parquet_options.allow_cast_unsigned_ints = true;

let mut builder = ParquetExecBuilder::new(file_scan_config)
.with_table_parquet_options(table_parquet_options)
.with_schema_adapter_factory(Arc::new(SparkSchemaAdapterFactory::new(
spark_cast_options,
spark_parquet_options,
)));

if let Some(filter) = cnf_data_filters {
Expand Down
9 changes: 5 additions & 4 deletions native/core/src/parquet/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,8 @@ use datafusion::datasource::listing::PartitionedFile;
use datafusion::datasource::physical_plan::parquet::ParquetExecBuilder;
use datafusion::datasource::physical_plan::FileScanConfig;
use datafusion::physical_plan::ExecutionPlan;
use datafusion_comet_spark_expr::{EvalMode, SparkCastOptions, SparkSchemaAdapterFactory};
use datafusion_comet_spark_expr::parquet_support::SparkParquetOptions;
use datafusion_comet_spark_expr::{EvalMode, SparkSchemaAdapterFactory};
use datafusion_common::config::TableParquetOptions;
use datafusion_execution::{SendableRecordBatchStream, TaskContext};
use futures::{poll, StreamExt};
Expand Down Expand Up @@ -679,13 +680,13 @@ pub unsafe extern "system" fn Java_org_apache_comet_parquet_Native_initRecordBat
table_parquet_options.global.pushdown_filters = true;
table_parquet_options.global.reorder_filters = true;

let mut spark_cast_options = SparkCastOptions::new(EvalMode::Legacy, "UTC", false);
spark_cast_options.allow_cast_unsigned_ints = true;
let mut spark_parquet_options = SparkParquetOptions::new(EvalMode::Legacy, "UTC", false);
spark_parquet_options.allow_cast_unsigned_ints = true;

let builder2 = ParquetExecBuilder::new(file_scan_config)
.with_table_parquet_options(table_parquet_options)
.with_schema_adapter_factory(Arc::new(SparkSchemaAdapterFactory::new(
spark_cast_options,
spark_parquet_options,
)));

//TODO: (ARROW NATIVE) - predicate pushdown??
Expand Down
2 changes: 2 additions & 0 deletions native/spark-expr/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,8 @@ pub use normalize_nan::NormalizeNaNAndZero;
mod variance;
pub use variance::Variance;
mod comet_scalar_funcs;
pub mod parquet_support; // TODO: Do `pub use` below to expose only what we need.

pub use cast::{spark_cast, Cast, SparkCastOptions};
pub use comet_scalar_funcs::create_comet_physical_fun;
pub use error::{SparkError, SparkResult};
Expand Down
Loading