Math Expressions#

Abs#

The following cases have no native implementation and always run in the JVM using Spark’s code-generated implementation (inside the Comet pipeline):

  • INTERVAL YEAR TO MONTH and INTERVAL DAY TO SECOND inputs

Greatest#

For applicable cases that are not selected for native execution automatically, Greatest is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline) by default. Set spark.comet.expression.Greatest.allowIncompatible=true to explicitly select Comet’s native implementation, which has the following differences from Spark:

  • Spark evaluates non-UTF8_BINARY collated string input under its collation, while Comet’s native greatest compares raw bytes

Least#

For applicable cases that are not selected for native execution automatically, Least is evaluated in the JVM using Spark’s own code-generated implementation (run inside the Comet pipeline) by default. Set spark.comet.expression.Least.allowIncompatible=true to explicitly select Comet’s native implementation, which has the following differences from Spark:

  • Spark evaluates non-UTF8_BINARY collated string input under its collation, while Comet’s native least compares raw bytes

Rand#

The following cases are not supported by Comet and always fall back to Spark, regardless of any allowIncompatible setting:

  • The seed argument must be a literal value

Randn#

The following cases are not supported by Comet and always fall back to Spark, regardless of any allowIncompatible setting:

  • The seed argument must be a literal value

Round#

The following cases have no native implementation and always run in the JVM using Spark’s code-generated implementation (inside the Comet pipeline):

  • Float and double inputs. Spark rounds them through a BigDecimal built from java.lang.Double.toString() rather than from the exact binary value, and that shortened decimal string can round differently than the value it came from

  • Negative-scale decimal inputs, which are only creatable with spark.sql.legacy.allowNegativeScaleOfDecimal=true