Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
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
Original file line number Diff line number Diff line change
Expand Up @@ -576,6 +576,36 @@ class VeloxSparkPlanExecApi extends SparkPlanExecApi with Logging {
VeloxHashExpressionTransformer(substraitExprName, exprs, original)
}

override def genFormatNumberTransformer(
substraitExprName: String,
children: Seq[ExpressionTransformer],
original: FormatNumber): ExpressionTransformer = {
if (!GlutenConfig.get.enableColumnarFormatNumber) {
throw new GlutenNotSupportException(
"Native format_number is disabled by spark.gluten.sql.columnar.formatNumber")
}
// Velox's format_number only supports an integer decimal-places second argument.
// Spark also accepts a Java DecimalFormat pattern STRING (e.g., '#,##0.00') which
// Velox does not implement -- fall back to vanilla Spark for that overload.
original.right.dataType match {
case StringType =>
throw new GlutenNotSupportException(
"format_number with string format pattern is not supported in native path")
case _ => // IntegerType -- supported
}
// DecimalType inputs require full-precision BigDecimal formatting which Velox does
// not support natively. Casting to Double would lose precision for values with >15
// significant digits. Fall back to Spark for DecimalType inputs.
original.left.dataType match {
case _: DecimalType =>
throw new GlutenNotSupportException(
"format_number with DecimalType input is not supported in native path " +
"(would lose precision when casting to Double)")
case _ => // Integer and floating-point types -- supported
}
GenericExpressionTransformer(substraitExprName, children, original)
}

/**
* Generate ShuffleDependency for ColumnarShuffleExchangeExec.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,9 +17,10 @@
package org.apache.gluten.execution

import org.apache.spark.SparkConf
import org.apache.spark.sql.catalyst.expressions.{Alias, Literal}
import org.apache.spark.sql.catalyst.expressions.{Alias, FormatNumber, Literal}
import org.apache.spark.sql.catalyst.optimizer.{ConstantFolding, NullPropagation}
import org.apache.spark.sql.classic.ClassicColumn
import org.apache.spark.sql.execution.ProjectExec
import org.apache.spark.sql.functions.col
import org.apache.spark.sql.types.StringType

Expand Down Expand Up @@ -678,4 +679,87 @@ class VeloxStringFunctionsSuite extends VeloxWholeStageTransformerSuite {
s"select l_orderkey, unbase64(base64(l_comment)) " +
s"from $LINEITEM_TABLE limit $LENGTH")(checkGlutenPlan[ProjectExecTransformer])
}

// format_number integration tests

private def formatNumberInProject(
p: org.apache.spark.sql.execution.SparkPlan): Boolean =
p.expressions.exists(_.exists(_.isInstanceOf[FormatNumber]))

// Asserts format_number WAS offloaded to Velox: a ProjectExecTransformer must carry the
// FormatNumber expression (proving native offload, not merely that some unrelated
// transformer exists in the plan).
private def assertFormatNumberOffloaded(df: org.apache.spark.sql.DataFrame): Unit = {
val plan = stripAQEPlan(df.queryExecution.executedPlan)
assert(
collectWithSubqueries(plan) {
case p: ProjectExecTransformer if formatNumberInProject(p) => p
}.nonEmpty,
"format_number should be offloaded to a ProjectExecTransformer")
}

// Asserts format_number was NOT offloaded to Velox: no ProjectExecTransformer may carry a
// FormatNumber expression, and a vanilla Spark ProjectExec must (proving genuine fallback
// rather than the expression merely being absent from the projection).
private def assertFormatNumberFallsBack(df: org.apache.spark.sql.DataFrame): Unit = {
val plan = stripAQEPlan(df.queryExecution.executedPlan)
assert(
collectWithSubqueries(plan) {
case p: ProjectExecTransformer if formatNumberInProject(p) => p
}.isEmpty,
"format_number must not be offloaded to a ProjectExecTransformer")
assert(
collectWithSubqueries(plan) { case p: ProjectExec if formatNumberInProject(p) => p }.nonEmpty,
"format_number should execute in a vanilla Spark ProjectExec")
}

test("format_number executes natively for integer input") {
runQueryAndCompare(
s"select l_orderkey, format_number(l_orderkey, 2) " +
s"from $LINEITEM_TABLE limit $LENGTH")(assertFormatNumberOffloaded)
}

test("format_number executes natively for double input") {
runQueryAndCompare(
s"select l_orderkey, format_number(CAST(l_orderkey AS DOUBLE), 4) " +
s"from $LINEITEM_TABLE limit $LENGTH")(assertFormatNumberOffloaded)
}

Comment on lines +716 to +727
test("format_number executes natively for float input") {
runQueryAndCompare(
s"select l_orderkey, format_number(CAST(l_orderkey AS FLOAT), 3) " +
s"from $LINEITEM_TABLE limit $LENGTH")(assertFormatNumberOffloaded)
}

test("format_number executes natively for bigint input") {
runQueryAndCompare(
s"select l_orderkey, format_number(CAST(l_orderkey AS BIGINT), 0) " +
s"from $LINEITEM_TABLE limit $LENGTH")(assertFormatNumberOffloaded)
}

test("format_number executes natively with zero decimal places") {
runQueryAndCompare(
s"select l_orderkey, format_number(l_orderkey, 0) " +
s"from $LINEITEM_TABLE limit $LENGTH")(assertFormatNumberOffloaded)
}

test("format_number falls back for string format pattern") {
runQueryAndCompare(
s"select l_orderkey, format_number(CAST(l_orderkey AS DOUBLE), '#,##0.00') " +
s"from $LINEITEM_TABLE limit $LENGTH")(assertFormatNumberFallsBack)
}

test("format_number falls back for DecimalType input") {
runQueryAndCompare(
s"select l_orderkey, format_number(CAST(l_orderkey AS DECIMAL(10,2)), 2) " +
s"from $LINEITEM_TABLE limit $LENGTH")(assertFormatNumberFallsBack)
}

test("format_number falls back when the feature is disabled") {
withSQLConf("spark.gluten.sql.columnar.formatNumber" -> "false") {
runQueryAndCompare(
s"select l_orderkey, format_number(l_orderkey, 2) " +
s"from $LINEITEM_TABLE limit $LENGTH")(assertFormatNumberFallsBack)
}
}
}
1 change: 1 addition & 0 deletions docs/Configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@ nav_order: 15
| spark.gluten.sql.columnar.filter | 🔄 Dynamic | true | Enable or disable columnar filter. |
| spark.gluten.sql.columnar.force.hashagg | 🔄 Dynamic | true | Whether to force to use gluten's hash agg for replacing vanilla spark's sort agg. |
| spark.gluten.sql.columnar.forceShuffledHashJoin | 🔄 Dynamic | true |
| spark.gluten.sql.columnar.formatNumber | 🔄 Dynamic | true | Enable or disable native execution of format_number(numeric, int) via Velox. Enabled by default. When enabled, integer and floating-point inputs are lowered to the Velox native implementation; the format-pattern STRING overload (e.g., '#,##0.00') and DecimalType inputs are unsupported and always fall back to Spark. |
| spark.gluten.sql.columnar.generate | 🔄 Dynamic | true |
| spark.gluten.sql.columnar.hashagg | 🔄 Dynamic | true | Enable or disable columnar hashagg. |
| spark.gluten.sql.columnar.hivetablescan | 🔄 Dynamic | true | Enable or disable columnar hivetablescan. |
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -532,6 +532,20 @@ trait SparkPlanExecApi {
GenericExpressionTransformer(substraitExprName, exprs, original)
}

/**
* Hook for native `format_number(numeric, int)` lowering. The default implementation throws
* [[GlutenNotSupportException]] so backends that do not support a native path cleanly fall back
* to vanilla Spark. Backends that do support the function (e.g., Velox) override this and may
* additionally enforce gating and type whitelisting.
*/
def genFormatNumberTransformer(
substraitExprName: String,
children: Seq[ExpressionTransformer],
original: FormatNumber): ExpressionTransformer = {
throw new GlutenNotSupportException(
"format_number native path is not supported by this backend")
}

/** Define backend-specific expression mappings. */
def extraExpressionMappings: Seq[Sig] = Seq.empty

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,8 @@ class GlutenConfig(conf: SQLConf) extends GlutenCoreConfig(conf) {

def enableColumnarWindowGroupLimit: Boolean = getConf(COLUMNAR_WINDOW_GROUP_LIMIT_ENABLED)

def enableColumnarFormatNumber: Boolean = getConf(COLUMNAR_FORMAT_NUMBER_ENABLED)

def enableAppendData: Boolean = getConf(COLUMNAR_APPEND_DATA_ENABLED)

def enableReplaceData: Boolean = getConf(COLUMNAR_REPLACE_DATA_ENABLED)
Expand Down Expand Up @@ -943,6 +945,16 @@ object GlutenConfig extends ConfigRegistry {
.booleanConf
.createWithDefault(true)

val COLUMNAR_FORMAT_NUMBER_ENABLED =
buildConf("spark.gluten.sql.columnar.formatNumber")
.doc(
"Enable or disable native execution of format_number(numeric, int) via Velox. " +
"Enabled by default. When enabled, integer and floating-point inputs are lowered " +
"to the Velox native implementation; the format-pattern STRING overload (e.g., " +
"'#,##0.00') and DecimalType inputs are unsupported and always fall back to Spark.")
.booleanConf
.createWithDefault(true)

val COLUMNAR_APPEND_DATA_ENABLED =
buildConf("spark.gluten.sql.columnar.appendData")
.doc("Enable or disable columnar v2 command append data.")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -939,6 +939,12 @@ object ExpressionConverter extends SQLConfHelper with Logging {
substraitExprName,
replaceWithExpressionTransformer0(errorMessage, attributeSeq, expressionsMap),
re)
case fn: FormatNumber =>
BackendsApiManager.getSparkPlanExecApiInstance.genFormatNumberTransformer(
substraitExprName,
fn.children.map(replaceWithExpressionTransformer0(_, attributeSeq, expressionsMap)),
fn
)
case expr =>
GenericExpressionTransformer(
substraitExprName,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,7 @@ object ExpressionMappings {
Sig[UnBase64](UNBASE64),
Sig[Base64](BASE64),
Sig[FormatString](FORMAT_STRING),
Sig[FormatNumber](FORMAT_NUMBER),

// URL functions
Sig[ParseUrl](PARSE_URL),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -143,6 +143,7 @@ object ExpressionNames {
final val BASE64 = "base64"
final val MASK = "mask"
final val FORMAT_STRING = "format_string"
final val FORMAT_NUMBER = "format_number"
final val LUHN_CHECK = "luhn_check"
final val TO_PRETTY_STRING = "to_pretty_string"

Expand Down
Loading