peter-toth commented on code in PR #59251:
URL: https://github.com/apache/spark/pull/59251#discussion_r4205060350
##########
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala:
##########
@@ -186,8 +186,15 @@ case class TransformExpression(
case None => throw QueryExecutionErrors.cannotEvaluateExpressionError(this)
}
+ // `resolvedFunction` is not a child, so a tree walk does not see it. For a
function with only
Review Comment:
Thanks. The `doGenCode` change is gone in 913f317a345. No output ordering
holds a transform now, so the `resolvedFunction` scaladoc stays as it is.
##########
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala:
##########
@@ -186,8 +186,15 @@ case class TransformExpression(
case None => throw QueryExecutionErrors.cannotEvaluateExpressionError(this)
}
+ // `resolvedFunction` is not a child, so a tree walk does not see it. For a
function with only
+ // `produceResult` it falls back to `eval` on the input row, which
whole-stage codegen does not
+ // provide (see `CodegenFallback.generate`). No whole-stage operator
generates code for a
+ // transform today.
Review Comment:
The `doGenCode` change is reverted in 913f317a345. A transform in
whole-stage codegen fails with INTERNAL_ERROR again, as on master. So no
`CollapseCodegenStages` change is needed.
##########
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala:
##########
@@ -186,8 +186,15 @@ case class TransformExpression(
case None => throw QueryExecutionErrors.cannotEvaluateExpressionError(this)
}
+ // `resolvedFunction` is not a child, so a tree walk does not see it. For a
function with only
+ // `produceResult` it falls back to `eval` on the input row, which
whole-stage codegen does not
+ // provide (see `CodegenFallback.generate`). No whole-stage operator
generates code for a
+ // transform today.
override protected def doGenCode(ctx: CodegenContext, ev: ExprCode):
ExprCode =
Review Comment:
Taken in 913f317a345. `LazyCodeGenOrdering` is now `LazyRowOrdering`, and it
builds its comparator with `RowOrdering.create`. That still happens in a
`@transient lazy val`, so the comparator is still built on the executor. A new
test in `GroupPartitionsExecSuite` serializes a `LazyRowOrdering` over a
transform and checks that it compares rows under `FALLBACK`.
##########
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala:
##########
@@ -186,8 +186,15 @@ case class TransformExpression(
case None => throw QueryExecutionErrors.cannotEvaluateExpressionError(this)
}
+ // `resolvedFunction` is not a child, so a tree walk does not see it. For a
function with only
+ // `produceResult` it falls back to `eval` on the input row, which
whole-stage codegen does not
+ // provide (see `CodegenFallback.generate`). No whole-stage operator
generates code for a
+ // transform today.
override protected def doGenCode(ctx: CodegenContext, ev: ExprCode):
ExprCode =
- throw QueryExecutionErrors.cannotGenerateCodeForExpressionError(this)
+ resolvedFunction match {
Review Comment:
Fixed in 913f317a345. The merge no longer compares rows by a transform, so
it never calls the function. The derived-ordering test now uses `days`. The
plan keeps the merge, over `[id]`.
##########
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala:
##########
@@ -186,8 +186,15 @@ case class TransformExpression(
case None => throw QueryExecutionErrors.cannotEvaluateExpressionError(this)
}
+ // `resolvedFunction` is not a child, so a tree walk does not see it. For a
function with only
+ // `produceResult` it falls back to `eval` on the input row, which
whole-stage codegen does not
+ // provide (see `CodegenFallback.generate`). No whole-stage operator
generates code for a
+ // transform today.
override protected def doGenCode(ctx: CodegenContext, ev: ExprCode):
ExprCode =
- throw QueryExecutionErrors.cannotGenerateCodeForExpressionError(this)
+ resolvedFunction match {
+ case Some(fn) => fn.genCode(ctx)
Review Comment:
Agreed, this predates the PR. After 913f317a345 the merge no longer
evaluates a transform, but the keyed shuffle still does. I filed SPARK-60027
for it.
##########
sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/expressions/TransformExpressionSuite.scala:
##########
@@ -112,4 +115,56 @@ class TransformExpressionSuite extends SparkFunSuite {
assert(!b12.hasSameReducedKeys(left))
assert(!b12.hasSameReducedKeys(b8), "nor do two unreduced ones")
}
+
+ test("SPARK-59995: the generated code computes what eval does") {
+ // `checkEvaluation` runs both the interpreted and the generated path.
Each function resolves to
+ // one of the calls `V2ExpressionUtils.resolveScalarFunction` builds.
+ val input = BoundReference(0, IntegerType, nullable = true)
+ Seq(
+ InvokeBucketFunction -> classOf[Invoke],
+ new StaticInvokeBucketFunction -> classOf[StaticInvoke],
+ ProduceResultBucketFunction -> classOf[ApplyFunctionExpression]
+ ).foreach { case (fn, call) =>
+ val transform = bucket(fn, input)
+ assert(call.isInstance(transform.resolvedFunction.get),
transform.resolvedFunction)
Review Comment:
Removed with the `doGenCode` change in 913f317a345.
##########
sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/expressions/TransformExpressionSuite.scala:
##########
@@ -112,4 +115,56 @@ class TransformExpressionSuite extends SparkFunSuite {
assert(!b12.hasSameReducedKeys(left))
assert(!b12.hasSameReducedKeys(b8), "nor do two unreduced ones")
}
+
+ test("SPARK-59995: the generated code computes what eval does") {
+ // `checkEvaluation` runs both the interpreted and the generated path.
Each function resolves to
+ // one of the calls `V2ExpressionUtils.resolveScalarFunction` builds.
+ val input = BoundReference(0, IntegerType, nullable = true)
+ Seq(
+ InvokeBucketFunction -> classOf[Invoke],
+ new StaticInvokeBucketFunction -> classOf[StaticInvoke],
+ ProduceResultBucketFunction -> classOf[ApplyFunctionExpression]
+ ).foreach { case (fn, call) =>
+ val transform = bucket(fn, input)
+ assert(call.isInstance(transform.resolvedFunction.get),
transform.resolvedFunction)
+ checkEvaluation(transform, 3, create_row(7))
Review Comment:
Removed with the `doGenCode` change in 913f317a345.
##########
sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/expressions/TransformExpressionSuite.scala:
##########
@@ -112,4 +115,56 @@ class TransformExpressionSuite extends SparkFunSuite {
assert(!b12.hasSameReducedKeys(left))
assert(!b12.hasSameReducedKeys(b8), "nor do two unreduced ones")
}
+
+ test("SPARK-59995: the generated code computes what eval does") {
+ // `checkEvaluation` runs both the interpreted and the generated path.
Each function resolves to
+ // one of the calls `V2ExpressionUtils.resolveScalarFunction` builds.
+ val input = BoundReference(0, IntegerType, nullable = true)
+ Seq(
+ InvokeBucketFunction -> classOf[Invoke],
+ new StaticInvokeBucketFunction -> classOf[StaticInvoke],
+ ProduceResultBucketFunction -> classOf[ApplyFunctionExpression]
+ ).foreach { case (fn, call) =>
+ val transform = bucket(fn, input)
+ assert(call.isInstance(transform.resolvedFunction.get),
transform.resolvedFunction)
+ checkEvaluation(transform, 3, create_row(7))
+ checkEvaluation(transform, 2, create_row(6))
+ }
+ }
+
+ test("SPARK-59995: a transform whose keys a join reduced generates no code")
{
+ // The call no longer computes such keys, so it must not run in the
generated path either.
+ val input = BoundReference(0, IntegerType, nullable = true)
Review Comment:
Removed with the `doGenCode` change in 913f317a345.
##########
sql/catalyst/src/test/scala/org/apache/spark/sql/catalyst/expressions/TransformExpressionSuite.scala:
##########
@@ -112,4 +115,56 @@ class TransformExpressionSuite extends SparkFunSuite {
assert(!b12.hasSameReducedKeys(left))
assert(!b12.hasSameReducedKeys(b8), "nor do two unreduced ones")
}
+
+ test("SPARK-59995: the generated code computes what eval does") {
+ // `checkEvaluation` runs both the interpreted and the generated path.
Each function resolves to
+ // one of the calls `V2ExpressionUtils.resolveScalarFunction` builds.
+ val input = BoundReference(0, IntegerType, nullable = true)
+ Seq(
+ InvokeBucketFunction -> classOf[Invoke],
+ new StaticInvokeBucketFunction -> classOf[StaticInvoke],
+ ProduceResultBucketFunction -> classOf[ApplyFunctionExpression]
+ ).foreach { case (fn, call) =>
+ val transform = bucket(fn, input)
+ assert(call.isInstance(transform.resolvedFunction.get),
transform.resolvedFunction)
+ checkEvaluation(transform, 3, create_row(7))
+ checkEvaluation(transform, 2, create_row(6))
+ }
+ }
+
+ test("SPARK-59995: a transform whose keys a join reduced generates no code")
{
+ // The call no longer computes such keys, so it must not run in the
generated path either.
+ val input = BoundReference(0, IntegerType, nullable = true)
+ val reduced = bucket(InvokeBucketFunction, input, 12)
+ .reducedTogetherWith(bucket(InvokeBucketFunction, input, 8))
+ checkError(
+ exception = intercept[SparkException](reduced.genCode(new
CodegenContext)),
+ condition = "INTERNAL_ERROR",
+ parameters = Map("message" -> s"Cannot generate code for expression:
$reduced"))
+ }
+}
+
+/** `bucket` over integers. The functions below differ only in how Spark calls
them. */
+private trait IntBucketFunction extends ScalarFunction[Int] {
Review Comment:
Removed with the `doGenCode` change in 913f317a345.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]