Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

feat: Support Ansi mode in abs function #500

Merged
merged 28 commits into from
Jun 11, 2024
Merged
Show file tree
Hide file tree
Changes from 14 commits
Commits
Show all changes
28 commits
Select commit Hold shift + click to select a range
fce78fa
change proto msg
planga82 May 30, 2024
1071eee
QueryPlanSerde with eval mode
planga82 May 30, 2024
750331d
Move eval mode
planga82 May 30, 2024
9f89b57
Add abs in planner
planga82 May 30, 2024
e6eda86
CometAbsFunc wrapper
planga82 May 30, 2024
0b37f8e
Add error management
planga82 May 30, 2024
d1e2099
Add tests
planga82 May 31, 2024
73e5513
Add license
planga82 May 31, 2024
cff5f29
spotless apply
planga82 May 31, 2024
9b3b4c8
format
planga82 May 31, 2024
708fffe
Merge remote-tracking branch 'refs/remotes/upstream/main' into bugfix…
planga82 May 31, 2024
f7df357
Fix clippy
planga82 Jun 1, 2024
76914b0
error msg for all spark versions
planga82 Jun 1, 2024
3b55ca2
Fix benches
planga82 Jun 1, 2024
aa92450
Merge upstream/main
planga82 Jun 3, 2024
ab28bf6
Use enum to ansi mode
planga82 Jun 3, 2024
0dda0b2
Fix format
planga82 Jun 3, 2024
1fc4f48
Add more tests
planga82 Jun 4, 2024
828ab3b
Merge remote-tracking branch 'refs/remotes/upstream/main' into bugfix…
planga82 Jun 4, 2024
fe2a003
Format
planga82 Jun 4, 2024
dc3f2a8
Merge remote-tracking branch 'refs/remotes/upstream/main' into bugfix…
planga82 Jun 4, 2024
3dff4bb
Refactor
planga82 Jun 5, 2024
19969d6
refactor
planga82 Jun 5, 2024
6fb873a
Merge remote-tracking branch 'refs/remotes/upstream/main' into bugfix…
planga82 Jun 6, 2024
bf64a24
Merge remote-tracking branch 'refs/remotes/upstream/main' into bugfix…
planga82 Jun 8, 2024
b4df447
merge upstream master
planga82 Jun 8, 2024
a72db13
fix merge
planga82 Jun 8, 2024
809052d
fix merge
planga82 Jun 8, 2024
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
2 changes: 1 addition & 1 deletion core/benches/cast_from_string.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@

use arrow_array::{builder::StringBuilder, RecordBatch};
use arrow_schema::{DataType, Field, Schema};
use comet::execution::datafusion::expressions::cast::{Cast, EvalMode};
use comet::execution::datafusion::expressions::{cast::Cast, EvalMode};
use criterion::{criterion_group, criterion_main, Criterion};
use datafusion_physical_expr::{expressions::Column, PhysicalExpr};
use std::sync::Arc;
Expand Down
2 changes: 1 addition & 1 deletion core/benches/cast_numeric.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@

use arrow_array::{builder::Int32Builder, RecordBatch};
use arrow_schema::{DataType, Field, Schema};
use comet::execution::datafusion::expressions::cast::{Cast, EvalMode};
use comet::execution::datafusion::expressions::{cast::Cast, EvalMode};
use criterion::{criterion_group, criterion_main, Criterion};
use datafusion_physical_expr::{expressions::Column, PhysicalExpr};
use std::sync::Arc;
Expand Down
3 changes: 3 additions & 0 deletions core/src/errors.rs
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,9 @@ pub enum CometError {
to_type: String,
},

#[error("[ARITHMETIC_OVERFLOW] {from_type} overflow. If necessary set \"spark.sql.ansi.enabled\" to \"false\" to bypass this error.")]
ArithmeticOverflow { from_type: String },

#[error(transparent)]
Arrow {
#[from]
Expand Down
87 changes: 87 additions & 0 deletions core/src/execution/datafusion/expressions/abs.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
// Licensed to the Apache Software Foundation (ASF) under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you under the Apache License, Version 2.0 (the
// "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.

use arrow::datatypes::DataType;
use arrow_schema::ArrowError;
use datafusion::logical_expr::{ColumnarValue, ScalarUDFImpl, Signature};
use datafusion_common::DataFusionError;
use datafusion_functions::math;
use std::{any::Any, sync::Arc};

use crate::errors::CometError;

use super::EvalMode;

fn arithmetic_overflow_error(from_type: &str) -> CometError {
CometError::ArithmeticOverflow {
from_type: from_type.to_string(),
}
}
Copy link
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This seems to duplicate the same function from negative.rs. Perhaps that one could be moved so that it can be reused here.


#[derive(Debug)]
pub struct CometAbsFunc {
inner_abs_func: Arc<dyn ScalarUDFImpl>,
eval_mode: EvalMode,
data_type_name: String,
}

impl CometAbsFunc {
pub fn new(eval_mode: EvalMode, data_type_name: String) -> Self {
Self {
inner_abs_func: math::abs().inner(),
eval_mode,
data_type_name,
}
}
}

impl ScalarUDFImpl for CometAbsFunc {
fn as_any(&self) -> &dyn Any {
self
}
fn name(&self) -> &str {
"abs"
}

fn signature(&self) -> &Signature {
self.inner_abs_func.signature()
}

fn return_type(&self, arg_types: &[DataType]) -> Result<DataType, DataFusionError> {
self.inner_abs_func.return_type(arg_types)
}

fn invoke(&self, args: &[ColumnarValue]) -> Result<ColumnarValue, DataFusionError> {
match self.inner_abs_func.invoke(args) {
Ok(result) => Ok(result),
Copy link
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Technically there is no need to match on an Ok result because there is already the catch all other handling.

Suggested change
Ok(result) => Ok(result),

Err(DataFusionError::ArrowError(ArrowError::ComputeError(msg), trace))
if msg.contains("overflow") =>
Copy link
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It would be nice if Arrow/DataFusion threw a specific overflow error so that we didn't have to look for a string within the error message, but I guess that isn't available.

Copy link
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I am going to file a feature request in DataFusion. I will post the link here later.

Copy link
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copy link
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Great! I understand that we have to wait for the next version of arror-rs to be released and integrate it here to be able to make the changes, right?

Copy link
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We can continue with this and do a follow up once the arrow-rs change is available. Or you can wait; the arrow-rs community is very responsive.

Copy link
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Probably better to do it in a follow-up PR, I think. Thanks!

Copy link
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It would help to log an issue to keep track

Copy link
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Apparently my PR in Arrow will not be available for 3 months because it is an API change 😞

Copy link
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

When we upgrade to the version of the arrow with the improved overflow eror reporing then the tests in this PR will fail (because we are looking for ComputeError but instead will get ArithmeticOverflow) so I don't think we need to file an issue

Copy link
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@parthchandra fyi ☝️

{
if self.eval_mode == EvalMode::Legacy {
Ok(args[0].clone())
} else {
let msg = arithmetic_overflow_error(&self.data_type_name).to_string();
Err(DataFusionError::ArrowError(
ArrowError::ComputeError(msg),
trace,
))
}
}
other => other,
}
}
}
9 changes: 2 additions & 7 deletions core/src/execution/datafusion/expressions/cast.rs
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,8 @@ use crate::{
},
};

use super::EvalMode;

static TIMESTAMP_FORMAT: Option<&str> = Some("%Y-%m-%d %H:%M:%S%.f");

static CAST_OPTIONS: CastOptions = CastOptions {
Expand All @@ -60,13 +62,6 @@ static CAST_OPTIONS: CastOptions = CastOptions {
.with_timestamp_format(TIMESTAMP_FORMAT),
};

#[derive(Debug, Hash, PartialEq, Clone, Copy)]
pub enum EvalMode {
Legacy,
Ansi,
Try,
}

#[derive(Debug, Hash)]
pub struct Cast {
pub child: Arc<dyn PhysicalExpr>,
Expand Down
8 changes: 8 additions & 0 deletions core/src/execution/datafusion/expressions/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ pub mod if_expr;
mod normalize_nan;
pub mod scalar_funcs;
pub use normalize_nan::NormalizeNaNAndZero;
pub mod abs;
pub mod avg;
pub mod avg_decimal;
pub mod bloom_filter_might_contain;
Expand All @@ -37,3 +38,10 @@ pub mod sum_decimal;
pub mod temporal;
mod utils;
pub mod variance;

#[derive(Debug, Hash, PartialEq, Clone, Copy)]
pub enum EvalMode {
Legacy,
Ansi,
Try,
}
36 changes: 23 additions & 13 deletions core/src/execution/datafusion/planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,6 @@ use datafusion::{
arrow::{compute::SortOptions, datatypes::SchemaRef},
common::DataFusionError,
execution::FunctionRegistry,
functions::math,
logical_expr::{
BuiltinScalarFunction, Operator as DataFusionOperator, ScalarFunctionDefinition,
},
Expand Down Expand Up @@ -52,6 +51,7 @@ use datafusion_common::{
tree_node::{Transformed, TransformedResult, TreeNode, TreeNodeRecursion, TreeNodeRewriter},
JoinType as DFJoinType, ScalarValue,
};
use datafusion_physical_expr::udf::ScalarUDF;
use itertools::Itertools;
use jni::objects::GlobalRef;
use num::{BigInt, ToPrimitive};
Expand All @@ -65,7 +65,7 @@ use crate::{
avg_decimal::AvgDecimal,
bitwise_not::BitwiseNotExpr,
bloom_filter_might_contain::BloomFilterMightContain,
cast::{Cast, EvalMode},
cast::Cast,
checkoverflow::CheckOverflow,
correlation::Correlation,
covariance::Covariance,
Expand Down Expand Up @@ -95,6 +95,8 @@ use crate::{
},
};

use super::expressions::{abs::CometAbsFunc, EvalMode};

// For clippy error on type_complexity.
type ExecResult<T> = Result<T, ExecutionError>;
type PhyAggResult = Result<Vec<Arc<dyn AggregateExpr>>, ExecutionError>;
Expand Down Expand Up @@ -149,6 +151,20 @@ impl PhysicalPlanner {
}
}

fn eval_mode_from_str(
eval_mode_str: &str,
allow_try: bool,
) -> Result<EvalMode, ExecutionError> {
match eval_mode_str {
"ANSI" => Ok(EvalMode::Ansi),
"LEGACY" => Ok(EvalMode::Legacy),
"TRY" if allow_try => Ok(EvalMode::Try),
other => Err(ExecutionError::GeneralError(format!(
"Invalid EvalMode: \"{other}\""
))),
}
}

/// Create a DataFusion physical expression from Spark physical expression
fn create_expr(
&self,
Expand Down Expand Up @@ -348,16 +364,7 @@ impl PhysicalPlanner {
let child = self.create_expr(expr.child.as_ref().unwrap(), input_schema)?;
let datatype = to_arrow_datatype(expr.datatype.as_ref().unwrap());
let timezone = expr.timezone.clone();
let eval_mode = match expr.eval_mode.as_str() {
"ANSI" => EvalMode::Ansi,
"TRY" => EvalMode::Try,
"LEGACY" => EvalMode::Legacy,
other => {
return Err(ExecutionError::GeneralError(format!(
"Invalid Cast EvalMode: \"{other}\""
)))
}
};
let eval_mode = Self::eval_mode_from_str(expr.eval_mode.as_str(), true)?;
Ok(Arc::new(Cast::new(child, datatype, eval_mode, timezone)))
}
ExprStruct::Hour(expr) => {
Expand Down Expand Up @@ -495,7 +502,10 @@ impl PhysicalPlanner {
let child = self.create_expr(expr.child.as_ref().unwrap(), input_schema.clone())?;
let return_type = child.data_type(&input_schema)?;
let args = vec![child];
let scalar_def = ScalarFunctionDefinition::UDF(math::abs());
let eval_mode = Self::eval_mode_from_str(expr.eval_mode.as_str(), false)?;
let comet_abs =
ScalarUDF::new_from_impl(CometAbsFunc::new(eval_mode, return_type.to_string()));
let scalar_def = ScalarFunctionDefinition::UDF(Arc::new(comet_abs));

let expr =
ScalarFunctionExpr::new("abs", scalar_def, args, return_type, None, false);
Expand Down
1 change: 1 addition & 0 deletions core/src/execution/proto/expr.proto
Original file line number Diff line number Diff line change
Expand Up @@ -473,6 +473,7 @@ message BitwiseNot {

message Abs {
Expr child = 1;
string eval_mode = 2;
}

message Subquery {
Expand Down
22 changes: 14 additions & 8 deletions spark/src/main/scala/org/apache/comet/serde/QueryPlanSerde.scala
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,13 @@ import org.apache.comet.shims.ShimQueryPlanSerde
* An utility object for query plan and expression serialization.
*/
object QueryPlanSerde extends Logging with ShimQueryPlanSerde with CometExprShim {

object ExecutionMode {
val ANSI = "ANSI"
val LEGACY = "LEGACY"
val TRY = "TRY"
}

def emitWarning(reason: String): Unit = {
logWarning(s"Comet native execution is disabled due to: $reason")
}
Expand Down Expand Up @@ -691,7 +698,7 @@ object QueryPlanSerde extends Logging with ShimQueryPlanSerde with CometExprShim
case Cast(child, dt, timeZoneId, evalMode) =>
val evalModeStr = if (evalMode.isInstanceOf[Boolean]) {
// Spark 3.2 & 3.3 has ansiEnabled boolean
if (evalMode.asInstanceOf[Boolean]) "ANSI" else "LEGACY"
if (evalMode.asInstanceOf[Boolean]) ExecutionMode.ANSI else ExecutionMode.LEGACY
} else {
// Spark 3.4+ has EvalMode enum with values LEGACY, ANSI, and TRY
evalMode.toString
Expand Down Expand Up @@ -1474,15 +1481,14 @@ object QueryPlanSerde extends Logging with ShimQueryPlanSerde with CometExprShim
None
}

case Abs(child, _) =>
case Abs(child, failOnErr) =>
Copy link
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

since we are already using failOnErr, can we simply use this boolean instead of evalmode struct? what are you thoughts?

Copy link
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hi, thanks for your review! The intention why it is done this way is that in the Rust code the ansi mode is always treated the same way whether the expression supports the three ansi modes or only two. If not, we have to do a different treatment for both cases.

val childExpr = exprToProtoInternal(child, inputs)
if (childExpr.isDefined) {
val abs =
ExprOuterClass.Abs
.newBuilder()
.setChild(childExpr.get)
.build()
Some(Expr.newBuilder().setAbs(abs).build())
val evalModeStr = if (failOnErr) ExecutionMode.ANSI else ExecutionMode.LEGACY
val absBuilder = ExprOuterClass.Abs.newBuilder()
absBuilder.setChild(childExpr.get)
absBuilder.setEvalMode(evalModeStr)
Some(Expr.newBuilder().setAbs(absBuilder).build())
} else {
withInfo(expr, child)
None
Expand Down
28 changes: 28 additions & 0 deletions spark/src/test/scala/org/apache/comet/CometExpressionSuite.scala
Original file line number Diff line number Diff line change
Expand Up @@ -850,6 +850,34 @@ class CometExpressionSuite extends CometTestBase with AdaptiveSparkPlanHelper {
}
}

test("abs Overflow ansi mode") {
val data: Seq[(Int, Int)] = Seq((Int.MaxValue, Int.MinValue))
Copy link
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can we have tests for all numerical values?

Copy link
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

done!

withSQLConf(
SQLConf.ANSI_ENABLED.key -> "true",
CometConf.COMET_ANSI_MODE_ENABLED.key -> "true") {
withParquetTable(data, "tbl") {
checkSparkMaybeThrows(sql("select abs(_1), abs(_2) from tbl")) match {
case (Some(sparkExc), Some(cometExc)) =>
val cometErrorPattern =
""".+[ARITHMETIC_OVERFLOW].+overflow. If necessary set "spark.sql.ansi.enabled" to "false" to bypass this error.""".r
val sparkErrorPattern = ".*integer overflow.*".r
assert(cometErrorPattern.findFirstIn(cometExc.getMessage).isDefined)
assert(sparkErrorPattern.findFirstIn(sparkExc.getMessage).isDefined)
case _ => fail("Exception should be thrown")
}
}
}
}

test("abs Overflow legacy mode") {
val data: Seq[(Int, Int)] = Seq((Int.MaxValue, Int.MinValue), (1, -1))
withSQLConf(SQLConf.ANSI_ENABLED.key -> "false") {
withParquetTable(data, "tbl") {
checkSparkAnswerAndOperator("select abs(_1), abs(_2) from tbl")
}
}
}

test("ceil and floor") {
Seq("true", "false").foreach { dictionary =>
withSQLConf(
Expand Down
Loading