dongjoon-hyun commented on code in PR #59251:
URL: https://github.com/apache/spark/pull/59251#discussion_r4197936286


##########
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:
   Nit: "a function with only `produceResult`" is narrower than when this falls 
back. `V2ExpressionUtils.findMethod` uses `getDeclaredMethod`, so the call is 
also an `ApplyFunctionExpression` when `invoke` is inherited or declared with 
other parameter classes; the PR description's "declares ... in its own class" 
is the accurate wording. Separately, the `resolvedFunction` scaladoc (L169-176) 
explains why its reduced-keys arm is unreachable only for partitionings and the 
write path ("every consumer of a reduced partitioning refuses it first"), while 
`doGenCode` now reaches it through an ordering. It is unreachable there because 
only `GroupPartitionsExec`'s output partitioning reports a reduced transform, 
and the only ordering built from partition keys is the scan's own.
   
   Could we say "when Spark calls the function through `produceResult`" here, 
and add the ordering case to that scaladoc?



##########
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:
   Minor: this checks the class of `resolvedFunction` and the values, but 
nothing checks that `doGenCode` emits the call's own code. If the `Some(fn)` 
branch fell back to `CodegenFallback.generate(this, ctx, ev)`, 
`checkEvaluation` here (`INPUT_ROW` "i") and the merge comparator in 
KeyGroupedPartitioningSuite (`INPUT_ROW` "a"/"b") would all still pass, so the 
main claim of the PR, that the `Invoke`/`StaticInvoke` code is generated 
instead of falling back to `eval`, is not pinned.
   
   Could we assert for the two `invoke` fixtures that the generated code calls 
`invoke(` directly and does not call `eval(`?



##########
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:
   Nit: this test never binds, evaluates or generates its child, since 
`doGenCode` throws before generating any child code, so it does not need a 
second `BoundReference`. Using `a`, as the SPARK-59121 test does for the same 
reduced shape (L101-104), removes the duplicate and checks "generates no code" 
more strictly: `a` is unevaluable, so a change that generated the child's code 
first would fail here.
   
   Could we use `a` here?



##########
sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala:
##########
@@ -6571,6 +6571,78 @@ class KeyGroupedPartitioningSuite
     }
   }
 
+  /** Two rows with `id` 1, each on its own split, so a join on `id` coalesces 
them. */
+  private def insertItemsInTwoYears(): Unit = sql(s"INSERT INTO 
testcat.ns.$items VALUES " +
+    "(1, 'aa', 10.0, cast('2021-01-01' as timestamp)), " +
+    "(1, 'ab', 11.0, cast('2022-01-01' as timestamp)), " +
+    "(2, 'bb', 20.0, cast('2021-01-01' as timestamp))")
+
+  /**
+   * Joins `items` with `purchases` on `joinCondition` and checks that the 
plan k-way merges over an
+   * ordering with a partition transform.
+   */
+  private def checkKWayMergeOverTransform(joinCondition: String): Unit = {
+    val df = sql(
+      s"""
+         |${selectWithMergeJoinHint("i", "p")}
+         |i.id, i.name, i.arrive_time
+         |FROM testcat.ns.$items i JOIN testcat.ns.$purchases p ON 
$joinCondition
+         |""".stripMargin)
+    checkAnswer(df, Seq(

Review Comment:
   Minor: neither new test can tell whether the merge orders rows by the value 
the transform's generated code computes, because `checkAnswer` ignores order. 
In the reported-ordering test, the only rows the merge compares are `(1, 'aa', 
2021)` and `(1, 'ab', 2022)`, and `name` decides before `years(arrive_time)`, 
so the transform's code is compiled but never run. In the key-derived test it 
runs, but the scan already sorts the splits by the full key, so a transform 
that returned a constant or a wrong value would still pass.
   
   Could the reported-ordering test make `id` and `name` tie, insert the later 
year first in a separate INSERT, and assert the merged order, e.g. with 
`merging.head.executeCollect()`?



##########
sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala:
##########
@@ -6571,6 +6571,78 @@ class KeyGroupedPartitioningSuite
     }
   }
 
+  /** Two rows with `id` 1, each on its own split, so a join on `id` coalesces 
them. */
+  private def insertItemsInTwoYears(): Unit = sql(s"INSERT INTO 
testcat.ns.$items VALUES " +
+    "(1, 'aa', 10.0, cast('2021-01-01' as timestamp)), " +
+    "(1, 'ab', 11.0, cast('2022-01-01' as timestamp)), " +
+    "(2, 'bb', 20.0, cast('2021-01-01' as timestamp))")
+
+  /**
+   * Joins `items` with `purchases` on `joinCondition` and checks that the 
plan k-way merges over an
+   * ordering with a partition transform.
+   */
+  private def checkKWayMergeOverTransform(joinCondition: String): Unit = {
+    val df = sql(
+      s"""
+         |${selectWithMergeJoinHint("i", "p")}
+         |i.id, i.name, i.arrive_time
+         |FROM testcat.ns.$items i JOIN testcat.ns.$purchases p ON 
$joinCondition
+         |""".stripMargin)
+    checkAnswer(df, Seq(
+      Row(1, "aa", Timestamp.valueOf("2021-01-01 00:00:00")),
+      Row(1, "ab", Timestamp.valueOf("2022-01-01 00:00:00")),
+      Row(2, "bb", Timestamp.valueOf("2021-01-01 00:00:00"))))
+    val merging = collectAllGroupPartitions(df.queryExecution.executedPlan)
+      .filter(_.enableSortedMerge)
+    assert(merging.length == 1, "expected one k-way merge")
+    
assert(merging.head.child.outputOrdering.exists(_.child.isInstanceOf[TransformExpression]),
+      "expected a partition transform in the merge's ordering")
+    assert(merging.head.execute().isInstanceOf[SortedMergeCoalescedRDD[_]])
+  }
+
+  test("SPARK-59995: k-way merge over a reported ordering with a partition 
transform") {
+    // The join on (id, name) needs the merge, since the key ordering on id 
alone is not enough.
+    // The merge's ordering is [id, name, years(arrive_time)], so generating 
its comparator
+    // generates code for the transform.
+    val itemOrdering = Array(
+      sort(FieldReference("id"), SortDirection.ASCENDING, 
NullOrdering.NULLS_FIRST),
+      sort(FieldReference("name"), SortDirection.ASCENDING, 
NullOrdering.NULLS_FIRST),
+      sort(years("arrive_time"), SortDirection.ASCENDING, 
NullOrdering.NULLS_FIRST))
+    createTable(items, itemsColumns, Array(identity("id")), itemOrdering)
+    insertItemsInTwoYears()
+    val namedPurchasesColumns = Array(
+      Column.create("item_id", LongType),
+      Column.create("name", StringType))
+    createTable(purchases, namedPurchasesColumns, Array(identity("item_id")))
+    sql(s"INSERT INTO testcat.ns.$purchases VALUES (1, 'aa'), (1, 'ab'), (2, 
'bb')")
+
+    withSQLConf(
+        SQLConf.REQUIRE_ALL_CLUSTER_KEYS_FOR_CO_PARTITION.key -> "false",
+        SQLConf.V2_BUCKETING_PRESERVE_ORDERING_ON_COALESCE_ENABLED.key -> 
"true") {
+      checkKWayMergeOverTransform("p.item_id = i.id AND p.name = i.name")
+    }
+  }
+
+  test("SPARK-59995: k-way merge over an ordering derived from a partition 
transform key") {
+    // The scan reports no ordering, so it derives [id, years(arrive_time)] 
from its keys. The join
+    // on id projects the keys to id. With preserveKeyOrderingOnCoalesce off, 
the coalesced
+    // partitions keep no ordering on id unless they are merged.
+    createTable(items, itemsColumns, Array(identity("id"), 
years("arrive_time")))

Review Comment:
   Minor: both new tests use `years`, which binds to `YearsFunction` with a 
magic `invoke`, so no end-to-end test runs a `produceResult`-only transform 
(the `ApplyFunctionExpression` route that the new comment in 
TransformExpression.scala singles out) through the merge's generated 
comparator; the catalyst test reaches that route only through projections. The 
suite already registers `UnboundBucketFunction`, whose `BucketFunction` has 
only `produceResult` and matches `items.id` (`LongType`).
   
   Could we add a variant with a `produceResult`-only key, e.g. 
`Array(identity("id"), bucket(4, "id"))`, so that the comparator evaluates it? 
Adding `bucket` to the first test's reported ordering would only compile it, 
since `name` decides first.



##########
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:
   **`resolvedFunction` is not a child, so the `CodegenFallback` it may hold is 
invisible to `CollapseCodegenStages`.** For a function called through 
`produceResult`, `fn` is an `ApplyFunctionExpression`, whose generated code is 
`((Expression) references[i]).eval(ctx.INPUT_ROW)`. The 
`CodegenFallback.generate` scaladoc (CodegenFallback.scala:47-50) says that a 
node falling back for a reason other than holding such a descendant "would have 
to be named in `CollapseCodegenStages` as well, since nothing else would tell 
it that the input row this code needs is not there", but 
`CollapseCodegenStages.supportCodegen` (WholeStageCodegenExec.scala:918-927) 
only walks children. I agree that no `CodegenSupport` operator holds a 
transform today, but the first one that does would fail worse than before: 
`consume` sets `ctx.INPUT_ROW = null` (WholeStageCodegenExec.scala:188), so 
`ProjectExec`/`FilterExec` would fail code generation with an NPE from the 
`code` interpolator, and where `INPUT_ROW` nam
 es another row (`HashAggregateExec`'s `unsafeRowBuffer`, a join's build row) 
the call would read the wrong row silently. Before this PR the same plan failed 
at codegen time with INTERNAL_ERROR.
   
   Could we add `case _: TransformExpression => false` to 
`CollapseCodegenStages.supportCodegen` (it changes nothing today) and mention 
`TransformExpression` in the `CodegenFallback.generate` scaladoc, instead of 
relying on this comment?



##########
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:
   **The merge comparator still has no interpreted fallback, so this fixes one 
expression rather than the mechanism that failed.** `LazyCodeGenOrdering` 
(GroupPartitionsExec.scala:806-810) calls `GenerateOrdering.generate` directly, 
while `SortExec.createSorter` uses `RowOrdering.create`, a 
`CodeGeneratorWithInterpretedFallback` that falls back to `InterpretedOrdering` 
on any NonFatal codegen error under the default 
`spark.sql.codegen.factoryMode`. So a transform whose call `eval` can run, but 
whose generated code cannot compile, still fails the task. For example, take a 
Java connector function declared as a package-private class with a public magic 
`invoke`: `Invoke.eval` works, since commons-lang3 
`MethodUtils.getMatchingAccessibleMethod` makes that method accessible, but the 
comparator casts to the class, which Janino should reject as inaccessible from 
the generated class's package. The same function works in a plain query or a 
`SortExec`, which log a warning and fall back. I unde
 rstand codegen was chosen because the comparator is on the hot path (#55116), 
and the sql/core tests run with `CODEGEN_ONLY` (SparkSessionBinder.scala:85), 
so they would keep exercising the generated path.
   
   Could `LazyCodeGenOrdering` use `RowOrdering.create(sortOrders, schema)`, or 
catch NonFatal itself, as a defense in depth?



##########
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:
   **The k-way merge is still planned when Spark cannot build this call, so 
such a query still fails at task time.** 
`GroupPartitionsExec.kWayMergeIsFeasible` (GroupPartitionsExec.scala:248-249) 
only checks `child.outputOrdering.nonEmpty && childIsSafeForKWayMerge`, and 
nothing at planning time checks that the ordering's transforms can be 
evaluated. This suite's own `days` is such a function: `DaysFunction` declares 
neither `invoke` nor `produceResult` in its own class, and 
`V2ExpressionUtils.findMethod` uses `getDeclaredMethod`, so building 
`resolvedFunction` here throws `SCALAR_FUNCTION_NOT_FULLY_IMPLEMENTED`. If the 
second new test partitions by `days("arrive_time")` instead of 
`years("arrive_time")`, the merge is planned over `[id, days(arrive_time)]` and 
the first comparison on the executor throws, while without the merge 
`EnsureRequirements` would add a `SortExec` on `id` and the query would return 
its rows. A bound function that is not a `ScalarFunction` fails the same way, wi
 th INTERNAL_ERROR from L196. This is not a regression, and the PR description 
scopes the fix to "as long as Spark can build the call", but `ScalarFunction`'s 
Javadoc says the magic method is resolved during query analysis, while here it 
is first resolved on the executors.
   
   Could `kWayMergeIsFeasible` refuse a merge whose ordering holds a transform 
without a buildable call, resolving it on the driver, and could we add a `days` 
test that checks the plan falls back to a `SortExec` and returns the right rows?



##########
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:
   Minor: the input is nullable, but only non-null rows are fed, so the 
generated null branches are never compared with `eval`. `Invoke` and 
`StaticInvoke` return null for a null primitive argument 
(`InvokeLike.needNullCheckForIndex`), while `ProduceResultBucketFunction` reads 
`input.getInt(1)` from `ApplyFunctionExpression`'s reused 
`SpecificInternalRow`, whose `getInt` ignores the null flag and returns 0. So 
the three fixtures are not interchangeable for a null input, despite "differ 
only in how Spark calls them" at L147.
   
   Could we add `create_row(null)` with a per-fixture expected value (null for 
the two `invoke` fixtures, 0 for `produceResult`), or make that fixture 
null-aware?



##########
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:
   Minor: this runs the call without the casts that 
`BoundFunction.inputTypes()` promises, now in generated code too. 
`resolvedFunction` (L177-182) passes the transform's raw children to 
`V2ExpressionUtils.resolveScalarFunction`, and nothing analyzes the resulting 
`Invoke`/`StaticInvoke`/`ApplyFunctionExpression`, so the cast promised in 
`BoundFunction`'s Javadoc ("Spark will cast input values to the required data 
types") never happens; only the write path coerces them (`TypeCoercionExecutor` 
in DistributionAndOrderingUtils.scala:86). With the test catalog, `bucket(4, 
c)` over an INT column fails with a `ClassCastException` in 
`ApplyFunctionExpression`'s `reusedRow`, whose slot is typed by the declared 
`LongType`, and `years(d)` over a DATE column, which `UnboundYearsFunction` 
accepts, reads the day count as microseconds. In the merge this only affects a 
tie-breaker that no parent relies on, but the keyed shuffle 
(ShuffleExchangeExec.scala:487) evaluates the same call. This predates 
 the PR, and the generated code behaves exactly like `eval`.
   
   Could `resolvedFunction` cast its arguments to `inputTypes()` as the write 
path does, or would you prefer a separate JIRA for it?



##########
sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala:
##########
@@ -6571,6 +6571,78 @@ class KeyGroupedPartitioningSuite
     }
   }
 
+  /** Two rows with `id` 1, each on its own split, so a join on `id` coalesces 
them. */
+  private def insertItemsInTwoYears(): Unit = sql(s"INSERT INTO 
testcat.ns.$items VALUES " +
+    "(1, 'aa', 10.0, cast('2021-01-01' as timestamp)), " +
+    "(1, 'ab', 11.0, cast('2022-01-01' as timestamp)), " +
+    "(2, 'bb', 20.0, cast('2021-01-01' as timestamp))")
+
+  /**
+   * Joins `items` with `purchases` on `joinCondition` and checks that the 
plan k-way merges over an
+   * ordering with a partition transform.
+   */
+  private def checkKWayMergeOverTransform(joinCondition: String): Unit = {
+    val df = sql(
+      s"""
+         |${selectWithMergeJoinHint("i", "p")}
+         |i.id, i.name, i.arrive_time
+         |FROM testcat.ns.$items i JOIN testcat.ns.$purchases p ON 
$joinCondition
+         |""".stripMargin)
+    checkAnswer(df, Seq(
+      Row(1, "aa", Timestamp.valueOf("2021-01-01 00:00:00")),
+      Row(1, "ab", Timestamp.valueOf("2022-01-01 00:00:00")),
+      Row(2, "bb", Timestamp.valueOf("2021-01-01 00:00:00"))))
+    val merging = collectAllGroupPartitions(df.queryExecution.executedPlan)
+      .filter(_.enableSortedMerge)
+    assert(merging.length == 1, "expected one k-way merge")
+    
assert(merging.head.child.outputOrdering.exists(_.child.isInstanceOf[TransformExpression]),

Review Comment:
   **No requirement can contain a transform, so the merge evaluates this key 
only as a tie-breaker nobody consumes.** `EnsureRequirements` enables the merge 
only when the merged ordering satisfies the parent's requirement, 
`SortOrder.orderingSatisfies` is prefix-based with semantic equality, and no 
`requiredChildOrdering` can hold a `TransformExpression`: SMJ, aggregate and 
window keys are query expressions, and the write path replaces transforms with 
their calls. So in these two tests the requirements are `[id, name]` and 
`[id]`, and `years(arrive_time)` only breaks ties after them. In the 
key-derived case every comparison within a coalesced group ties on `id`, so 
each one calls the connector function twice to compute a value that is constant 
per split. Cutting `kWayMergeOrdering` (GroupPartitionsExec.scala:275-276) and 
the merge branch of `outputOrdering` (GroupPartitionsExec.scala:374-377) before 
the first sort key that holds a transform would leave every merge decision 
unchanged,
  avoid those calls, and also avoid the failure for a transform Spark cannot 
call (see the comment on TransformExpression.scala:194).
   
   Could we consider that, either instead of generating the transform's code 
here or in addition to 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)
+      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:
   Nit: these fixtures do not override `canonicalName()`, whose default returns 
a new random UUID on every call (BoundFunction.java:104-110), so two transforms 
over the same fixture object get different `functionId`s, and 
`isSameFunction`/`hasSameReducedKeys` treat them as different functions. 
Nothing here depends on it today, but the fixtures are top-level, so other 
tests in this package can pick them up, and the reduced-keys test's message 
embeds a random UUID. `NamedFunction` above has a stable canonical name for the 
same reason.
   
   Could we add `override def canonicalName(): String = "test.bucket"` here?



-- 
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]

Reply via email to