Repository navigation
Conversation
| transform.children.size == 1 && isReference(transform.children.head) | ||
| // TransformExpression.collectLeaves() only returns column references, not literals. | ||
| // We need exactly one column reference per transform. | ||
| transform.collectLeaves().size == 1 |
There was a problem hiding this comment.
[P1] This widens support from transforms with one direct reference child to any transform whose collectLeaves() returns one leaf. Because Transform arguments can themselves be nested Transforms, V2ExpressionUtils can materialize shapes like outer(years(k)) and outer(days(k)). The new gate accepts both, keyPositions later maps them only by leaf k, and TransformExpression.isSameFunction compares only the outer function name plus literals. That can make storage partitionings with different nested child semantics look compatible and let SPJ skip a required shuffle, which can drop join matches. Keep the old direct-reference constraint, or compare the full non-literal child semantics before admitting these transforms.
[ 🤖 posted by Codex on behalf of sunchao 🤖 ]
There was a problem hiding this comment.
Thanks @sunchao for the review. Good point. I didn't think about nested transform. This is indeed widening the transforms, as collectLeaves() go through all the children in nested expressions.
Like you mentioned, adding direct-reference constraints here would be enough. I'll make that change by getting non literal children, and use same existing check (size and reference). With that we also don't have to then worry about TransformExpression.isSameFunction comparing only outer function.
With regards to adding support to nested transform for SPJ, very interesting idea, but we can take that as separate PR.
| private def extractParameters(expr: TransformExpression): ReducibleParameters = { | ||
| import scala.jdk.CollectionConverters._ | ||
| val values = expr.literalChildren.map { | ||
| case Literal(value, _) => value.asInstanceOf[AnyRef] |
There was a problem hiding this comment.
[P1] extractParameters forwards raw Catalyst Literal.value objects into ReducibleParameters. For string literals, Spark stores UTF8String, while the new public API documents string parameters and exposes getString() as a java.lang.String cast. A connector implementing the documented string-parameter case will get ClassCastException from getString(0) or be forced to depend on Spark internals. Convert literal values to connector-facing external values by dataType before constructing ReducibleParameters.
[ 🤖 posted by Codex on behalf of sunchao 🤖 ]
There was a problem hiding this comment.
I agree, due UTF8String, we'll get ClassCastException. I think same case for Decimal type.
I'll make the change to handle these, and possibly a little more generic way for future proofing.
| val thisParams = extractParameters(thisExpr) | ||
| val otherParams = extractParameters(otherExpr) | ||
|
|
||
| val res = if (!thisParams.isEmpty && !otherParams.isEmpty) { |
There was a problem hiding this comment.
[P2] The new generalized reducer API accepts ReducibleParameters on both sides, and this file even models zero-literal transforms as ReducibleParameters([]). But this dispatch only invokes the generalized overload when both sides are non-empty. Mixed cases such as parameterized-vs-zero-parameter transforms instead fall back to reducer(otherFunction), so connectors that correctly implement the new generalized overload for those cases are never invoked; the default legacy path can even throw UnsupportedOperationException. If mixed arity is meant to be unsupported, the new API/docs should say that explicitly. Otherwise this should dispatch through the generalized overload whenever either side wants that path.
[ 🤖 posted by Codex on behalf of sunchao 🤖 ]
There was a problem hiding this comment.
@sunchao Good point. I was trying to fallback to existing behavior, but hadn't considered parameterized-vs-zero-parameter case, which it should follow the new generalized path. With regards making mixed arity "unsupported", currently I can't think of an valid example, but I also can't think of not to support it either. So, I suppose, we can let the implementor reducer decide to have it or not.
I will make the change that will dispatch through the generalized overloaded reducer().
| ReducibleFunction<?, ?> otherFunction, | ||
| ReducibleParameters otherParams) { | ||
| // Default: try old Int-based API for backward compatibility | ||
| if (thisParams.count() == 1 && otherParams.count() == 1) { |
There was a problem hiding this comment.
i think it makes more sense to have it the other way. the old one should call the new (more generic one). we can have the temporary code on spark side:
val thisParams = extractParameters(thisExpr)
val otherParams = extractParameters(otherExpr)
if (singleIntParam(thisParams) && singleIntParam(otherParams)) {
Option(thisFunction.reducer(thisParams.getInt(0), otherFunction, otherParams.getInt(0)))
.orElse(Option(thisFunction.reducer(thisParams, otherFunction, otherParams)))
} else if (!thisParams.isEmpty && !otherParams.isEmpty) {
Option(thisFunction.reducer(thisParams, otherFunction, otherParams))
} else {
Option(thisFunction.reducer(otherFunction))
}
imo, its not much uglier than this code :)
In practice, Iceberg is the only implementation of this so we can hopefully remove this ugly wrapper from spark side.
There was a problem hiding this comment.
@szehon-ho Yes, this make sense. My initial idea was to not call "deprecated" method, and when time to remove it would be only from 1 single file (ReducibleFunction), but the approach is not natural (or norm), where deprecated code calls the newer method.
When there's time to remove the method, we can remove from both, would not be a big deal, and can't be missed too.
I'll make this change along with @sunchao mixed arity case.
| * <li>custom_transform(col, "param") → ReducibleParameters(["param"])</li> | ||
| * </ul> | ||
| * | ||
| * @since 4.0.0 |
There was a problem hiding this comment.
@szehon-ho oh good eye! I worked on this last year :) Will update to reflect current version.
|
@sunchao @szehon-ho Addressed all comments, PTAL Additionally, I'm now handling all "null" or unimplemented exception (UOE) from cc @peter-toth |
peter-toth
left a comment
There was a problem hiding this comment.
Summary
Adds SPJ support for the truncate partition transform by generalizing the reducer API: literal parameters now live inline in TransformExpression.children, are surfaced to connectors through a new ReducibleParameters container, and a generalized ReducibleFunction.reducer(ReducibleParameters, ReducibleFunction, ReducibleParameters) overload replaces the bucket-only reducer(int, …, int) signature.
Prior state and problem. Bucket's numBuckets was carried as a sidecar numBucketsOpt: Option[Int] field on TransformExpression, outside the expression tree. The SPJ reducer dispatch in TransformExpression.reducer only knew how to fan out two cases — (Some, Some) → reducer(int, …, int) or anything else → reducer(otherFunction) — so any parameterized transform other than bucket (e.g. truncate(col, N)) was structurally invisible to the reducer dispatch and always shuffled. The write side was fixed in SPARK-40295; the read/join side never was. This PR supersedes the earlier exploration in #49211 by generalizing the API rather than special-casing transforms one at a time.
Design approach. Drop the numBucketsOpt sidecar and let literals participate in children. Then expose those literals to connectors as a typed parameter container, route the SPJ reducer dispatch through a single generalized overload, and keep the deprecated int-API alive (with default throw UOE) so existing connectors like Iceberg 1.10.0 continue to work. The dispatch tries the deprecated signature first for the single-int case (Iceberg shape), falls back to the new overload on UnsupportedOperationException, and goes straight to the new overload for everything else.
Key design decisions.
- Literals are now inline
Expressionchildren, not a typed field.TransformExpression.collectLeaves()is overridden to strip them so the existingpartitioning.scala/EnsureRequirementsconsumers still see exactly one leaf attribute per transform. ReducibleParametersis a brand-new public API rather than reusingV2Literal/Literal; it exposes typed getters (getInt,getLong,getString,getDouble,getFloat,get) but deliberately hides the sourceDataType.- Spark's own
BucketFunctionwas switched to override only the new API. Backward compat for connectors that override only the deprecated method is covered by a dedicatedLegacyBucketFunctiontest fixture. extractParametersconvertsUTF8String → java.lang.StringandDecimal → java.math.BigDecimal, but otherwise passes Catalyst-internal values straight through.
Implementation sketch. New public API in sql/catalyst/.../functions/ (ReducibleParameters.java, generalized ReducibleFunction.reducer overload). Catalyst plumbing in TransformExpression.scala (new literalChildren, extractParameters, tryReduce, dispatch logic; overrides for collectLeaves and canonicalized). V2-to-Catalyst resolution simplified in V2ExpressionUtils (the BucketTransform case is dropped; bucket now resolves through the generic NamedTransform path). Physical layer in partitioning.scala updates KeyedPartitioning.supportsExpressions to count non-literal children and KeyedShuffleSpec.createPartitioning to preserve literals when rewriting children. DistributionAndOrderingUtils.resolveTransformExpression simplified. Test fixtures in transformFunctions.scala add TruncateFunction.reducer via the new API, plus IntegerTruncateFunction and LegacyBucketFunction. New tests in KeyGroupedPartitioningSuite.
Behavioral changes worth calling out.
truncate(col, N)now plansKeyedPartitioningwhere it previously plannedUnknownPartitioning. The existing testnon-clustered distribution: V2 function with multiple argswas renamed toclustered distribution: …to reflect this.EnsureRequirementsSuiteexprA–exprDflipped fromLiteral(1..4)toAttributeReferences because the new code distinguishes literals from refs in transform children. Tests usingbucket(N, exprA)now exercise a[Literal(N), AttributeReference]shape rather than[Literal(N), Literal(1)]. Worth a sanity check that no existing assertion relied on the prior Literal shape.
General
- @sunchao (Codex, May 15) and @szehon-ho (May 18) already raised the substantive issues with the earlier commit. The author has addressed: the
collectLeaves()widening onKeyedPartitioning.supportsExpressions(now useschildren.filterNot(_.isInstanceOf[Literal])directly), the mixed-arity dispatch (now routes through the generalized overload), the dispatch ordering preference (deprecated-first, new fallback per szehon's suggestion), and the@sinceversion onReducibleParameters. Nice turnaround on all of those. - The fix for @sunchao's
extractParametersP1 is only partial. The current code convertsUTF8String → StringandDecimal → BigDecimal, but the original critique was structural: convert bydataTyperather than special-casing each type. As written,BinaryType(byte[]),CalendarIntervalType,ArrayType/MapTypeCatalyst values, etc. still flow through the catch-allvalue.asInstanceOf[AnyRef]and a connector implementing the documented "any type" claim onReducibleParameterswill hit Catalyst-internal classes. The principled fix is to drive the conversion off the literal'sDataType(e.g., via existingCatalystTypeConverters.convertToScala) — or, taken further, to surface V2 literals directly through the API; see inline #6. - The PR description claim "New overload
ReducibleFunction.reducer(ReducibleParameters, ...)with a default that delegates to the deprecated single-int signature for backward compatibility" is out of date with the code. The new default throwsUnsupportedOperationException; it'sTransformExpression.reducer's Spark-side dispatch that tries the deprecated signature first. Worth updating the description so reviewers don't get a wrong mental model of the fallback mechanism. - The
collectLeavesoverride onTransformExpression(inline #3 below) is the kind of contract divergence that's better fixed at the abstraction layer it crosses. I've put up #56088 as a follow-up that migrates the SPJ call sites inpartitioning.scalaandEnsureRequirements.scalaoff_.collectLeaves()— toExpression.referenceswhere set semantics suffices, and a newExpression.collectAttributes(): Seq[Attribute]where positionalSeqis required. Pure refactor against today'smaster. Once #56088 lands, this PR can rebase on top and drop the override here entirely; the migrated call sites already filter literals correctly.
Suggested improvements
ReducibleParametersis missinggetBigDecimal.extractParametersconvertsDecimal → java.math.BigDecimalon the wire, but the container exposes no decimal getter, so connectors with a decimal parameter must use the untypedget(index)and cast manually — exactly the type-safety the wrapper was added to prevent. (Inline #6 below proposes a structural alternative that obviates this gap entirely.) [sql/catalyst/src/main/java/org/apache/spark/sql/connector/catalog/functions/ReducibleParameters.java:114]canonicalizedoverride is dead code and the doc justification doesn't hold. SPJ usesisSameFunction, notcanonicalized/semanticEquals, so the stated motivation ("bucket(4, tableA.id) and bucket(4, tableB.id) semantically equal") is moot — andAttributeReference.canonicalizedkeepsexprId, so the override wouldn't deliver that equality even if SPJ did read it. The override otherwise duplicates the default path. Remove. [sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala:93]collectLeavesoverride breaks the universalTreeNodecontract — drop on rebase after #56088. The follow-up migrates the SPJ call sites off_.collectLeaves()so the override here can simply be removed; see the General section. [sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala:111]UnboundTruncateFunction.bind()claimsBinaryTypeandLongTypesupport, but neither branch actually works.BinaryType → TruncateFunctionmismatches(StringType, IntegerType)inputTypes()and agetUTF8Stringread;LongType → IntegerTruncateFunctionmismatches(IntegerType, IntegerType)inputTypes()and agetIntread. Plus theIntegerTruncateFunctiondoc claim about "incompatible partition structures" is wrong. Drop the unused branches, or split per type with coverage. [sql/core/src/test/scala/org/apache/spark/sql/connector/catalog/functions/transformFunctions.scala:293]isSameFunctioncompares literals by value but ignores positional structure. The check tolerates literals at different positions inchildrenand ignores nested-transform children entirely (compares only the outer function name + literal values), sobucket(4, days(col))andbucket(4, hours(col))would be reported as the same function. A structural per-position compare closes both gaps. [sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala:75]- Structural alternative — reuse
V2Literalinstead of introducingReducibleParameters.V2Literalis already public and already used everywhere else in the connector API; passing it through the reducer signature preservesdataType()(no Catalyst-internal-types leak), makes inline #1 moot (no typed getters to maintain), and shrinksextractParametersto a Catalyst-Literal → V2-Literal conversion. Worth weighing before locking the new public class in. [sql/catalyst/src/main/java/org/apache/spark/sql/connector/catalog/functions/ReducibleParameters.java:40]
| */ | ||
| public float getFloat(int index) { | ||
| return (Float) values.get(index); | ||
| } |
There was a problem hiding this comment.
extractParameters in TransformExpression.scala:171 converts Decimal → java.math.BigDecimal on the way in, but this class exposes no getBigDecimal accessor. A connector with a decimal parameter therefore has to fall back to the untyped Object get(int index) and cast by hand — exactly the type-safety the wrapper was introduced to provide. Either add a getBigDecimal getter alongside the others, or remove the Decimal special-case in extractParameters and document the supported types explicitly. Concretely:
/**
* Get parameter at index as BigDecimal.
* @throws ClassCastException if parameter is not a BigDecimal
* @throws IndexOutOfBoundsException if index is invalid
*/
public java.math.BigDecimal getBigDecimal(int index) {
return (java.math.BigDecimal) values.get(index);
}Same gap exists for any other type the Spark side might convert (binary, interval, etc.) — see the ## General note about driving extractParameters off DataType. Inline #6 proposes a structural alternative that would make this whole class of issue moot.
| * This is crucial for Storage Partitioned Joins - we need bucket(4, tableA.id) and bucket(4, | ||
| * tableB.id) to be semantically equal so SPJ can be triggered. | ||
| */ | ||
| override lazy val canonicalized: Expression = { |
There was a problem hiding this comment.
The doc claim that motivates this override doesn't hold, and the override doesn't add anything over the inherited default. Two pieces:
The motivation doesn't hold. The doc says "we need bucket(4, tableA.id) and bucket(4, tableB.id) to be semantically equal so SPJ can be triggered". SPJ doesn't use canonicalized/semanticEquals on TransformExpression though — it dispatches through isSameFunction/isCompatible (see KeyedShuffleSpec.isExpressionCompatible at partitioning.scala:1082). And even if it did, this override wouldn't deliver the claim: AttributeReference.canonicalized (namedExpressions.scala:324) normalizes the name to "none" but keeps the original exprId, so two different attributes with the same name still canonicalize to different forms. bucket(4, tableA.id) and bucket(4, tableB.id) aren't made equal by either the override or the default.
The body is equivalent to the default. Literal is a LeafExpression, so Literal.canonicalized returns the literal unchanged — the per-child case l: Literal => l branch produces the same result as l.canonicalized. The override therefore matches withCanonicalizedChildren exactly, except it skips the Canonicalize.execute post-pass that the default routes through (a no-op for TransformExpression today, but a future change to Canonicalize would silently fail to apply here).
Recommended: remove the override entirely.
There was a problem hiding this comment.
Good catch, yes, not needed. Was meant to remove it.
| * | ||
| */ | ||
| override def collectLeaves(): Seq[Expression] = { | ||
| children.flatMap { |
There was a problem hiding this comment.
This override breaks the universal TreeNode.collectLeaves contract — the base method (TreeNode.scala:315) returns "every node in the tree where children.isEmpty", but this returns column references only, hiding the literal leaves. Functionally correct for current SPJ call sites (partitioning.scala:492/498/982/1037, EnsureRequirements.scala:89/412/417/780), but a quiet semantic divergence — any future caller of collectLeaves on a TransformExpression will silently get wrong answers.
The cleaner shape is to migrate the SPJ call sites off _.collectLeaves() and use the right helper at each:
Expression.referenceswhere set semantics suffice (existence/size checks, single-attribute lookup).- A new
Expression.collectAttributes(): Seq[Attribute]where a positionalSeqis required (RowOrdering.create,reorder,attributes.zip(clustering)) —AttributeSetdeduplicates and its ordering isLinkedHashSet-implementation-detail rather than contract.
I've sent #56088 as the follow-up doing exactly that. Pure refactor against today's master (current TransformExpression has no literal children, so the migrated call sites produce identical results). Once it's merged, this override can be dropped on rebase, and 55885's representation change (literals as inline children) lands cleanly through the already-migrated call sites.
There was a problem hiding this comment.
Yes, exactly this. I was getting away with the context of TransformExpression as it always meant for Attributes, but with respect to TreeNode.collectLeaves contract, it was indeed breaking.
I was also trying handle by filter Literals during usages, but meant modifying all the callers.
collectAttributes() is exactly what I needed, and will use it once #56088 land on main branch.
Thanks for the PR @peter-toth
There was a problem hiding this comment.
Actually, now I think we don't even need a new Expression.collectAttributes(), but we can use Expression.references() in SPJ planning. Adjusted the PR.
There was a problem hiding this comment.
I've merged #56088 so collectLeaves() override is no longer needed.
| override def bind(inputType: StructType): BoundFunction = { | ||
| if (inputType.size == 2) { | ||
| inputType.head.dataType match { | ||
| case StringType | BinaryType => TruncateFunction |
There was a problem hiding this comment.
This dispatch claims BinaryType and LongType support, but neither branch actually works:
StringType | BinaryType => TruncateFunction—TruncateFunction.inputTypes()declares(StringType, IntegerType)andproduceResultcallsinput.getUTF8String(0). ABinaryTyperow (rawbyte[]underneath) mismatchesinputTypes()at bind-time, or — if the framework lets it through — fails at runtime on thegetUTF8Stringread.IntegerType | LongType => IntegerTruncateFunction—IntegerTruncateFunction.inputTypes()declares(IntegerType, IntegerType)andproduceResultcallsinput.getInt(0). Same kind of mismatch on aLongTyperow, plus itsresultType()isIntegerTypeso even if you cast input down, the column type is wrong.
None of the phantom branches are exercised by tests in this PR — every truncate test uses StringType and dispatches to TruncateFunction correctly. So the broken branches add API surface that connector authors might trust without coverage to defend it.
Adjacent issue: the doc on IntegerTruncateFunction ("different integer truncate widths produce incompatible partition structures") is wrong: truncate(v, W1) is reducible to truncate(v, W2) whenever W2 % W1 == 0 (snap to a coarser width), and when neither divides the other both reduce to truncate(_, lcm(W1, W2)). Same structural argument as bucket's GCD case.
Three options, in order of cleanliness:
- Drop the unused branches and
IntegerTruncateFunctionfrom this PR — onlyStringTypetruncate is part of the SPJ story being delivered. - Split per type (
IntegerTruncate/LongTruncate/BinaryTruncate) with the rightinputTypes()/resultType()/produceResultand add at least one test per branch. - Make each existing object polymorphic on the input type (declare wider
inputTypes(), dispatch inproduceResult) and add tests.
There was a problem hiding this comment.
@peter-toth fixed Removed unsupported datatype in the test. Also I converted IntegerTruncateFunction with proper reducible method, along with new unit test for it.
| // Compare literal arguments to ensure transforms with different parameters | ||
| // (e.g., bucket(32, col) vs bucket(16, col), truncate(col, 2) vs truncate(col, 4)) | ||
| // are not considered the same | ||
| val otherLiterals = other.literalChildren |
There was a problem hiding this comment.
This compares literals by value but ignores positional structure of children. Two consequences:
-
Position-agnostic. Two transforms with the same function and the same literal values but different interleavings of literals/non-literals are reported as the same. In practice the V2 function catalog pins down argument order per function name, so this is unlikely to manifest, but the doc claims "same semantics as
other" — the check is weaker than that without the catalog assumption. -
Nested-transform children.
V2ExpressionUtils.toCatalystrecurses throughtoCatalystTransformOpt, so children can include nestedTransformExpressions (sunchao's earlier P1 raised this forsupportsExpressions).bucket(4, days(col))andbucket(4, hours(col))both have outer functionbucketandliteralChildren = [Literal(4)], so this returnstrueeven though the inner transforms differ.KeyedPartitioning.supportsExpressionsrejects nested-transform shapes from SPJ planning, so this doesn't bite SPJ today, butisSameFunctionis used elsewhere (isCompatible, suite assertions) where the false positive would leak.
A structural per-position walk closes both gaps and lets literalChildren shrink to just the extractParameters helper:
def isSameFunction(other: TransformExpression): Boolean = {
function.canonicalName() == other.function.canonicalName() &&
children.length == other.children.length &&
children.zip(other.children).forall {
case (a: Literal, b: Literal) => a == b
case (a: TransformExpression, b: TransformExpression) => a.isSameFunction(b)
case (_: Literal, _) | (_, _: Literal) => false
case (_: TransformExpression, _) | (_, _: TransformExpression) => false
case _ => true // both refs/attrs — column identity intentionally ignored
}
}There was a problem hiding this comment.
@sunchao raised this issue too, and I stuck with the current limitation of SPJ with only 1 Reference.
But I like the idea of generalizing it to support nested cases as well.
I'll try to follow the suggestion and make the change.
Thanks @peter-toth for suggestion.
There was a problem hiding this comment.
@peter-toth
Updated isSameFunction method which is recursive:
- outer function name + per-position children;
- recurses into nested transforms;
- ignores column identity
So, nowbucket(4, years(c))is not the same asbucket(4, days(c)).
| * @since 5.0.0 | ||
| */ | ||
| @Evolving | ||
| public class ReducibleParameters { |
There was a problem hiding this comment.
Structural alternative — reuse V2Literal instead of introducing ReducibleParameters. Surfacing this even though it's late in the cycle, because it's the kind of design choice worth weighing before locking in a new public class.
The proposal: drop ReducibleParameters entirely and have the generalized reducer take org.apache.spark.sql.connector.expressions.Literal (the V2 literal type already used everywhere else in the connector API):
default Reducer<I, O> reducer(
org.apache.spark.sql.connector.expressions.Literal<?>[] thisParams,
ReducibleFunction<?, ?> otherFunction,
org.apache.spark.sql.connector.expressions.Literal<?>[] otherParams) {
throw new UnsupportedOperationException();
}Connector use:
@Override
public Reducer<UTF8String, UTF8String> reducer(
Literal<?>[] thisParams,
ReducibleFunction<?, ?> otherFunc,
Literal<?>[] otherParams) {
if (otherFunc != TruncateFunction) return null;
int thisWidth = (Integer) thisParams[0].value();
int otherWidth = (Integer) otherParams[0].value();
...
}What this fixes simultaneously:
- Inline Removed reference to incubation in README.md. #1's gap evaporates. No typed getters to maintain — connectors call
.value()and dispatch on.dataType(). No more "we forgotgetBigDecimal" / "we'll needgetCalendarInterval" / etc. - The
extractParameterspartial-conversion concern (General-section bullet Removed reference to incubation in Spark user docs. #2) goes away.V2LiteralcarriesdataType()alongside the value, so the connector interprets it correctly regardless of whether it'sString,BigDecimal,byte[],CalendarInterval, etc. No Catalyst-internal types leaking, no type-by-type special-casing —extractParametersshrinks to a Catalyst-Literal → V2-Literal conversion (essentially the inverse of whatV2ExpressionUtils.toCatalystalready does for V2 literals). - One fewer public class to learn / maintain / stabilize.
V2Literalis already public, already@Evolving→ stable, already used by every connector that authors V2 transforms (Expressions.literal(N),Expressions.bucket(N, col)). Receiving them back through the reducer API is symmetric round-trip with how transforms are constructed — V2 connectors never have to cross the Catalyst boundary in this public surface.
Trade-offs being honest about:
- Slightly more verbose at the connector use site (
(Integer) params[0].value()vsparams.getInt(0)). Real but small. - Java generics + arrays awkwardness;
Literal<?>[]produces unchecked-warning ceremony.List<Literal<?>>orLiteral<?>...varargs are cleaner alternatives. - The deprecated int-API → new-API dispatch shim still needed (single int →
Literal<Integer>).
cc @sunchao @szehon-ho — would value your read on this trade-off, given your existing reviews of ReducibleParameters. Should the API rebase on V2Literal, or stay with the new class as drafted?
### What changes were proposed in this pull request? - Add `Expression.collectAttributes(): Seq[Attribute]` for cases that need a positional `Seq` (preserving order and duplicates). - Migrate SPJ-related uses of `_.collectLeaves()` in `partitioning.scala` and `EnsureRequirements.scala` to: - `Expression.references` / `AttributeSet.fromAttributeSets(...)` where set semantics suffice (existence/size checks, single-attribute lookup). - The new `Expression.collectAttributes()` where positional `Seq[Attribute]` is required (`RowOrdering.create`, `reorder`, `attributes.zip(clustering)`). - Drop the now-redundant `.map(_.asInstanceOf[Attribute])` cast at the `EnsureRequirements` site that previously typed the result of `collectLeaves`. ### Why are the changes needed? `TreeNode.collectLeaves()` returns every node in the tree where `children.isEmpty` — including `Literal`s. SPJ planning has always wanted attributes only, but with the existing partition expression layout (`TransformExpression.children = [col]`, parameters carried in a sidecar `numBucketsOpt: Option[Int]` field), the difference didn't surface in practice. Follow-up work that puts literal parameters directly into `TransformExpression.children` (e.g. `bucket(Literal(numBuckets), col)`, `truncate(col, Literal(width))`) — see SPARK-50593 / apache#55885 — would otherwise force `TransformExpression` to override `collectLeaves` to filter literals, breaking the universal `TreeNode.collectLeaves` contract for one expression type. Fixing the abstraction layer by using the right helper at each call site is cleaner than overriding `collectLeaves` inside `TransformExpression`. ### Does this PR introduce _any_ user-facing change? No. Pure refactor — for current partition expression shapes (no literal children) `_.collectLeaves()`, `_.references`, and `_.collectAttributes()` produce identical results. ### How was this patch tested? Covered by existing test suites that exercise the migrated call sites: `KeyGroupedPartitioningSuite`, `EnsureRequirementsSuite`, `ProjectedOrderingAndPartitioningSuite`. ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude Opus 4.7
### What changes were proposed in this pull request? - Document `AttributeSet`'s iteration-order contract on the class scaladoc: iteration via `iterator` / `foreach` / `flatMap` returns elements in insertion order (driven by the underlying `LinkedHashSet`). `toSeq` is called out as the explicit exception — it sorts by `(name, exprId.id)` for codegen stability (SPARK-18394). - Migrate seven SPJ-related uses of `_.collectLeaves()` in `partitioning.scala` and `EnsureRequirements.scala` to `_.references` / `AttributeSet.fromAttributeSets(...)`. Drops the now-redundant `.map(_.asInstanceOf[Attribute])` cast at the `EnsureRequirements:89` site. ### Why are the changes needed? `TreeNode.collectLeaves()` returns every node in the tree where `children.isEmpty`, including `Literal`s. SPJ planning has always wanted attributes only, but with the existing partition expression layout (`TransformExpression.children = [col]`, parameters carried in a sidecar `numBucketsOpt: Option[Int]` field), the difference didn't surface. Follow-up work (e.g. SPARK-50593 / apache#55885) that puts literal parameters directly into `TransformExpression.children` (`bucket(Literal(numBuckets), col)`, `truncate(col, Literal(width))`) would otherwise force `TransformExpression` to override `collectLeaves` to filter literals, breaking the universal `TreeNode.collectLeaves` contract for one expression type. `Expression.references` already returns attributes only (filtering literals and other non-attribute leaves), and its insertion-ordered iteration is exactly what positional binding (`RowOrdering.create`, `reorder`, `attributes.zip(clustering)`) requires. The per-partition-expression single-column rule (enforced by `KeyedPartitioning.supportsExpressions`) ensures within-expression dedup never matters here. Documenting the iteration-order contract lets these call sites rely on the order without implicit dependency on implementation detail. ### Does this PR introduce _any_ user-facing change? No. ### How was this patch tested? Covered by existing test suites that exercise the migrated call sites: `KeyGroupedPartitioningSuite`, `EnsureRequirementsSuite`, `ProjectedOrderingAndPartitioningSuite`. ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude Opus 4.7
### What changes were proposed in this pull request? - Document `AttributeSet`'s iteration-order contract on the class scaladoc: iteration via `iterator` / `foreach` / `flatMap` returns elements in insertion order (driven by the underlying `LinkedHashSet`). `toSeq` is called out as the explicit exception — it sorts by `(name, exprId.id)` for codegen stability (SPARK-18394). - Migrate seven SPJ-related uses of `_.collectLeaves()` in `partitioning.scala` and `EnsureRequirements.scala` to `_.references` / `AttributeSet.fromAttributeSets(...)`. Drops the now-redundant `.map(_.asInstanceOf[Attribute])` cast at the `EnsureRequirements:89` site. ### Why are the changes needed? `TreeNode.collectLeaves()` returns every node in the tree where `children.isEmpty`, including `Literal`s. SPJ planning has always wanted attributes only, but with the existing partition expression layout (`TransformExpression.children = [col]`, parameters carried in a sidecar `numBucketsOpt: Option[Int]` field), the difference didn't surface. Follow-up work (e.g. SPARK-50593 / apache#55885) that puts literal parameters directly into `TransformExpression.children` (`bucket(Literal(numBuckets), col)`, `truncate(col, Literal(width))`) would otherwise force `TransformExpression` to override `collectLeaves` to filter literals, breaking the universal `TreeNode.collectLeaves` contract for one expression type. `Expression.references` already returns attributes only (filtering literals and other non-attribute leaves), and its insertion-ordered iteration is exactly what positional binding (`RowOrdering.create`, `reorder`, `attributes.zip(clustering)`) requires. The per-partition-expression single-column rule (enforced by `KeyedPartitioning.supportsExpressions`) ensures within-expression dedup never matters here. Documenting the iteration-order contract lets these call sites rely on the order without implicit dependency on implementation detail. ### Does this PR introduce _any_ user-facing change? No. ### How was this patch tested? Covered by existing test suites that exercise the migrated call sites: `KeyGroupedPartitioningSuite`, `EnsureRequirementsSuite`, `ProjectedOrderingAndPartitioningSuite`. ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude Opus 4.7
### What changes were proposed in this pull request? - Document `AttributeSet`'s iteration-order contract on the class scaladoc: iteration via `iterator` / `foreach` / `flatMap` returns elements in insertion order (driven by the underlying `LinkedHashSet`). `toSeq` is called out as the explicit exception — it sorts by `(name, exprId.id)` for codegen stability (SPARK-18394). - Migrate seven SPJ-related uses of `_.collectLeaves()` in `partitioning.scala` and `EnsureRequirements.scala` to `_.references` / `AttributeSet.fromAttributeSets(...)`. Drops the now-redundant `.map(_.asInstanceOf[Attribute])` cast at the `EnsureRequirements:89` site. - Update `EnsureRequirementsSuite` synthetic fixtures (`exprA..D`) from `Literal(1..4)` to `AttributeReference`s. The literals were stand-ins for columns; under the migration's `_.references`-based attribute extraction, literal children produce empty `AttributeSet`s and trip the planner's "exactly one attribute per partition expression" assertions. Real partitionings can't reach those assertions with literal-only transforms because `KeyedPartitioning.supportsExpressions`'s `isReference` check filters them out — so this is a fixture-only update that preserves test intent. ### Why are the changes needed? `TreeNode.collectLeaves()` returns every node in the tree where `children.isEmpty`, including `Literal`s. SPJ planning has always wanted attributes only, but with the existing partition expression layout (`TransformExpression.children = [col]`, parameters carried in a sidecar `numBucketsOpt: Option[Int]` field), the difference didn't surface. Follow-up work (e.g. SPARK-50593 / apache#55885) that puts literal parameters directly into `TransformExpression.children` (`bucket(Literal(numBuckets), col)`, `truncate(col, Literal(width))`) would otherwise force `TransformExpression` to override `collectLeaves` to filter literals, breaking the universal `TreeNode.collectLeaves` contract for one expression type. `Expression.references` already returns attributes only (filtering literals and other non-attribute leaves), and its insertion-ordered iteration is exactly what positional binding (`RowOrdering.create`, `reorder`, `attributes.zip(clustering)`) requires. The per-partition-expression single-column rule (enforced by `KeyedPartitioning.supportsExpressions`) ensures within-expression dedup never matters here. Documenting the iteration-order contract lets these call sites rely on the order without implicit dependency on implementation detail. ### Does this PR introduce _any_ user-facing change? No. ### How was this patch tested? Covered by existing test suites that exercise the migrated call sites: `KeyGroupedPartitioningSuite`, `EnsureRequirementsSuite`, `ProjectedOrderingAndPartitioningSuite`. ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude Opus 4.7
[SPARK-57038][SQL] Use `Expression.references` in SPJ planning ### What changes were proposed in this pull request? - Document `AttributeSet`'s iteration-order contract on the class scaladoc: iteration via `iterator` / `foreach` / `flatMap` returns elements in insertion order (driven by the underlying `LinkedHashSet`). `toSeq` is called out as the explicit exception — it sorts by `(name, exprId.id)` for codegen stability (SPARK-18394). - Migrate seven SPJ-related uses of `_.collectLeaves()` in `partitioning.scala` and `EnsureRequirements.scala` to `_.references` / `AttributeSet.fromAttributeSets(...)`. Drops the now-redundant `.map(_.asInstanceOf[Attribute])` cast at the `EnsureRequirements:89` site. - Update `EnsureRequirementsSuite` synthetic fixtures (`exprA..D`) from `Literal(1..4)` to `AttributeReference`s. The literals were stand-ins for columns; under the migration's `_.references`-based attribute extraction, literal children produce empty `AttributeSet`s and trip the planner's "exactly one attribute per partition expression" assertions. Real partitionings can't reach those assertions with literal-only transforms because `KeyedPartitioning.supportsExpressions`'s `isReference` check filters them out — so this is a fixture-only update that preserves test intent. ### Why are the changes needed? `TreeNode.collectLeaves()` returns every node in the tree where `children.isEmpty`, including `Literal`s. SPJ planning has always wanted attributes only, but with the existing partition expression layout (`TransformExpression.children = [col]`, parameters carried in a sidecar `numBucketsOpt: Option[Int]` field), the difference didn't surface. Follow-up work (e.g. SPARK-50593 / #55885) that puts literal parameters directly into `TransformExpression.children` (`bucket(Literal(numBuckets), col)`, `truncate(col, Literal(width))`) would otherwise force `TransformExpression` to override `collectLeaves` to filter literals, breaking the universal `TreeNode.collectLeaves` contract for one expression type. `Expression.references` already returns attributes only (filtering literals and other non-attribute leaves), and its insertion-ordered iteration is exactly what positional binding (`RowOrdering.create`, `reorder`, `attributes.zip(clustering)`) requires. The per-partition-expression single-column rule (enforced by `KeyedPartitioning.supportsExpressions`) ensures within-expression dedup never matters here. Documenting the iteration-order contract lets these call sites rely on the order without implicit dependency on implementation detail. ### Does this PR introduce _any_ user-facing change? No. ### How was this patch tested? Covered by existing test suites that exercise the migrated call sites: `KeyGroupedPartitioningSuite`, `EnsureRequirementsSuite`, `ProjectedOrderingAndPartitioningSuite`. ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude Opus 4.7 Closes #56088 from peter-toth/SPARK-57038-revisit-collectleaves. Authored-by: Peter Toth <peter.toth@gmail.com> Signed-off-by: Peter Toth <peter.toth@gmail.com>
[SPARK-57038][SQL] Use `Expression.references` in SPJ planning ### What changes were proposed in this pull request? - Document `AttributeSet`'s iteration-order contract on the class scaladoc: iteration via `iterator` / `foreach` / `flatMap` returns elements in insertion order (driven by the underlying `LinkedHashSet`). `toSeq` is called out as the explicit exception — it sorts by `(name, exprId.id)` for codegen stability (SPARK-18394). - Migrate seven SPJ-related uses of `_.collectLeaves()` in `partitioning.scala` and `EnsureRequirements.scala` to `_.references` / `AttributeSet.fromAttributeSets(...)`. Drops the now-redundant `.map(_.asInstanceOf[Attribute])` cast at the `EnsureRequirements:89` site. - Update `EnsureRequirementsSuite` synthetic fixtures (`exprA..D`) from `Literal(1..4)` to `AttributeReference`s. The literals were stand-ins for columns; under the migration's `_.references`-based attribute extraction, literal children produce empty `AttributeSet`s and trip the planner's "exactly one attribute per partition expression" assertions. Real partitionings can't reach those assertions with literal-only transforms because `KeyedPartitioning.supportsExpressions`'s `isReference` check filters them out — so this is a fixture-only update that preserves test intent. ### Why are the changes needed? `TreeNode.collectLeaves()` returns every node in the tree where `children.isEmpty`, including `Literal`s. SPJ planning has always wanted attributes only, but with the existing partition expression layout (`TransformExpression.children = [col]`, parameters carried in a sidecar `numBucketsOpt: Option[Int]` field), the difference didn't surface. Follow-up work (e.g. SPARK-50593 / #55885) that puts literal parameters directly into `TransformExpression.children` (`bucket(Literal(numBuckets), col)`, `truncate(col, Literal(width))`) would otherwise force `TransformExpression` to override `collectLeaves` to filter literals, breaking the universal `TreeNode.collectLeaves` contract for one expression type. `Expression.references` already returns attributes only (filtering literals and other non-attribute leaves), and its insertion-ordered iteration is exactly what positional binding (`RowOrdering.create`, `reorder`, `attributes.zip(clustering)`) requires. The per-partition-expression single-column rule (enforced by `KeyedPartitioning.supportsExpressions`) ensures within-expression dedup never matters here. Documenting the iteration-order contract lets these call sites rely on the order without implicit dependency on implementation detail. ### Does this PR introduce _any_ user-facing change? No. ### How was this patch tested? Covered by existing test suites that exercise the migrated call sites: `KeyGroupedPartitioningSuite`, `EnsureRequirementsSuite`, `ProjectedOrderingAndPartitioningSuite`. ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude Opus 4.7 Closes #56088 from peter-toth/SPARK-57038-revisit-collectleaves. Authored-by: Peter Toth <peter.toth@gmail.com> Signed-off-by: Peter Toth <peter.toth@gmail.com> (cherry picked from commit 408a115) Signed-off-by: Peter Toth <peter.toth@gmail.com>
…cibleParameters container
… Joins by generalizing parameter handling
…cibleParameters backward compatibility
86945cd to
721f576
Compare
|
@peter-toth @sunchao I've addressed most of the feedbacks except one (see below)
One item I haven't addressed yet is |
| // Use Expression.references (AttributeSet), which already filters out non-attribute | ||
| // leaves (e.g., Literal parameters introduced by parameterized transforms such as | ||
| // `bucket(numBuckets, col)` or `truncate(col, width)`). | ||
| transform.references.size == 1 |
There was a problem hiding this comment.
First, apologies if our earlier comments were unclear -- our suggestion to use references was specifically about consumption sites (see comment #55885 (comment)), not about widening this gate. We had not asked for the SPJ gate to be widened, and the recursive isSameFunction change does not require it.
The new commit bundles three semi-independent changes here and in TransformExpression:
- Recursive
isSameFunctioninTransformExpression-- the truncate fix. ✓ - Widening
isSupportedTransformtotransform.references.size == 1-- now admits nested transforms likebucket(4, years(col)), and alsobucket(4, cast(col)),bucket(4, col + 1), etc. - New
hasOnlyReferenceArgsgate onKeyedShuffleSpec.canCreatePartitioning, plus thenonLiteralChildrenSameguard insideTransformExpression.reducer-- both exist only because (2) admits shapes where the prior code was unsound.
(2) is a coverage expansion, not a fix for any existing bug. We'd like this PR to ship the truncate fix without the coverage expansion:
- Restore the previous gate that limited
isSupportedTransformto "single non-literal child that is a column reference." - Drop
hasOnlyReferenceArgs-- it exists only to narrow back what the widened gate would otherwise let through. - Drop the
nonLiteralChildrenSameguard insidereducer. It patches a real latent bug (without it,bucket(4, years(c))vsbucket(2, days(c))would silently reduce viagcd), but that bug is unreachable once the strict gate is restored: every nested-transform shape that reaches the reducer has already been ruled out by the gate itself.
Keep the recursive isSameFunction machinery -- it's the right shape for the truncate fix (position-agnostic literal compare) and is well-documented; future widening can build on it.
There was a problem hiding this comment.
Thanks @peter-toth
Reverted.
Yep I agree, for nested transforms support for SPJ, we can discuss separately.
| childrenMatch(other)(_ == _) | ||
|
|
||
| @scala.annotation.tailrec | ||
| private def isColumnRef(e: Expression): Boolean = e match { |
There was a problem hiding this comment.
This duplicates KeyedPartitioning.isReference in partitioning.scala -- same shape, same semantics. Conditional on inline #1 above (revert the widening, remove hasOnlyReferenceArgs), this method becomes unused and can be deleted. Otherwise please consolidate -- either call KeyedPartitioning.isReference from here, or hoist the helper into a shared utility used by both sites.
There was a problem hiding this comment.
Moved into companion objects for reusability.
| private def extractParameters(expr: TransformExpression): ReducibleParameters = { | ||
| import scala.jdk.CollectionConverters._ | ||
| val values = expr.literalChildren.map { | ||
| case Literal(value, dt) => CatalystTypeConverters.convertToScala(value, dt) |
There was a problem hiding this comment.
CatalystTypeConverters.convertToScala returns config-dependent runtime types for DateType and TimestampType:
DateType→java.time.LocalDateifspark.sql.datetime.java8API.enabledis true, otherwisejava.sql.Date.TimestampType→java.time.Instantif the same flag is true, otherwisejava.sql.Timestamp.
So the runtime class of a date/timestamp value inside ReducibleParameters is a function of the executor's SQLConf rather than of the column's DataType. A connector implementing reducer(ReducibleParameters, ReducibleFunction, ReducibleParameters) cannot just write params.get(0).asInstanceOf[LocalDate] -- it must handle both forms, or risk a ClassCastException on whichever form it didn't anticipate.
This is an additional point in favor of the V2Literal-array alternative discussed at #55885 (comment): V2 Literal<T> already carries the internal representation along with the DataType. Connectors receive a stable, self-describing (value, dataType) pair and decide for themselves how to interpret it -- which is the right contract for a public reducer API.
There was a problem hiding this comment.
Sounds good. I have Literal[] version ready go too. I'm doing some testing. Just waiting for @sunchao 's feedback.
There was a problem hiding this comment.
Thanks @peter-toth, I'm on board with reusing V2Literal over introducing ReducibleParameters. A few reasons it's the better long-term contract, and how I'd address the two weaknesses you called out:
Why V2Literal:
- It's already the public,
@Evolvingliteral type connectors use to construct transforms (Expressions.bucket(N, col),Expressions.literal(N)), so receiving them back throughreduceris a symmetric round-trip rather than a new concept to learn. - It's self-describing:
value()+dataType()lets the connector interpret the value itself. That structurally closes the issues we kept hitting withReducibleParameters— the missinggetBigDecimal, and the config-dependentDate/Timestampruntime types fromconvertToScala. WithV2Literalcarrying the internal representation alongside theDataType,extractParameterscollapses to a Catalyst-Literal-> V2-Literalconversion (the inverse of whatV2ExpressionUtils.toCatalystalready does), with no per-type special-casing. - Type-future-proofing: a new
DataTypeneeds zero API changes here, vs.ReducibleParametersneeding a new typed getter per type. One fewer public class to stabilize and maintain.
On the perceived weaknesses:
-
Complex/nested types leaking Catalyst internals. In practice transform parameters are scalars (
numBuckets, truncatewidth), so this doesn't bite. I'd rather make that explicit than leave it ambiguous: let's document on the reducer API that parameters are scalar literals, and haveextractParametersreject non-scalar literal children rather than silently forwarding anArrayData/MapData/InternalRow. That keeps the contract honest without expanding scope. (Note this caveat would apply toReducibleParameterstoo, only worse — it'd need a new getter and still leak.) -
Array-generics ceremony (
Literal<?>[]). Agree it's awkward. Let's use varargs (Literal<?>... thisParams/Literal<?>... otherParams) — it avoids the unchecked-warning noise, reads cleanly at the connector call site, and the int-API -> new-API shim just passesExpressions.literal(numBuckets).List<Literal<?>>would also work but varargs is the lightest touch.
The slightly-more-verbose call site ((Integer) params[0].value() vs params.getInt(0)) is a fair trade — it's the same idiom connectors already use against InternalRow, and Iceberg is migrating off the deprecated int-API regardless, so there's no existing implementation that loses out.
cc @sunchao — feel free to chime in if you want, since you reviewed ReducibleParameters too.
There was a problem hiding this comment.
@szehon-ho Thanks for the input, Just a note, I like the idea of varargs, but both Scala and Java doesn't support in the middle of the argument. Seems like we're stuck with Literal<?>[].
But I like the idea of not exposing the internal complex type/scalar literals.
There was a problem hiding this comment.
+1 on using V2 Literal<?>[] too, given that it can:
- Preserves both value() and dataType() instead of storing untyped Objects.
- Avoids maintaining separate getters like getInt(), getString(), and missing future getters.
- Handles typed nulls and new Spark data types without changing the public API.
- Avoids Spark converting values into configuration-dependent Java types.
- Reuses the existing connector literal abstraction instead of introducing another public class.
| @@ -4183,4 +4188,287 @@ class KeyGroupedPartitioningSuite extends DistributionAndOrderingSuiteBase with | |||
| } | |||
| } | |||
| } | |||
|
|
|||
| // === SPARK-50593: Truncate SPJ support tests === | |||
There was a problem hiding this comment.
Please drop this comment, the ticket prefixes identify the purpose of the tests.
|
Sorry for the delay. Will take another look! |
| thisFunction.reducer(numBuckets, otherFunction, otherNumBuckets) | ||
| case _ => thisFunction.reducer(otherFunction) | ||
| otherExpr: TransformExpression): Option[Reducer[_, _]] = { | ||
| val thisParams = extractParameters(thisExpr) |
There was a problem hiding this comment.
[P1] Please preserve the column/literal argument layout before invoking a reducer. The strict supportsExpressions gate only requires one direct column reference; it does not require that reference to occupy the same child position on both sides. extractParameters then removes the literal positions, and the latest commit removed the positional nonLiteralChildrenSame guard.
This is reachable for a public V2 function with two same-typed arguments. For example, f(id, 2) and f(4, store_id) both pass the gate. The generalized reducer sees only [2] and [4] and can report compatibility. I reproduced this on the current head with the integer-truncate reducer: both transforms were admitted, isCompatible returned true, and matching join value 3 mapped to reduced key 0 on the left versus key 3 on the right. SPJ can therefore skip a required shuffle and lose the match.
Please restore a positional-shape check before reducer dispatch: literal values may differ, but literal and column-reference slots must align on both sides.
Posted by Codex on behalf of @sunchao.
| } else if (isSingleInt(thisParams) && isSingleInt(otherParams)) { | ||
| // Try deprecated int-API first for legacy connectors (e.g. Iceberg 1.10); | ||
| // the first attempt is silent because we have a fallback. Only the fallback warns. | ||
| Try(Option(thisFunction.reducer( |
There was a problem hiding this comment.
[P2] Try(Option(oldReducer)).orElse(genericReducer) only evaluates the generalized overload when the deprecated call throws. A documented null return becomes Success(None), so Try.orElse returns it without executing the fallback.
I reproduced this on the current head with a function implementing both overloads: the deprecated overload returned null, the generalized overload returned a valid reducer, and reducers(...).isDefined was still false. A connector retaining the deprecated method while implementing broader compatibility through the replacement API will therefore miss a valid reducer and add an unnecessary shuffle.
Please inspect the inner Option and invoke the generalized overload after either a failed deprecated call or None. A regression test should cover legacy-null/new-success.
Posted by Codex on behalf of @sunchao.
…he reducer API
Carry reducible-function parameters as V2 Literal[] (value + DataType) instead of the
ReducibleParameters wrapper, and delete that class.
- ReducibleFunction.reducer(Literal[], ReducibleFunction, Literal[]) replaces the
ReducibleParameters overload; the deprecated reducer(int, ..., int) is kept for legacy
connectors (e.g. Iceberg 1.10).
- TransformExpression.extractParameters wraps each literal child as LiteralValue(value, dataType)
-- internal value + DataType, no CatalystTypeConverters conversion -- so connectors interpret
values via dataType() (resolves the config-dependent Date/Timestamp concern).
- Reducer dispatch hardening:
- sameArgumentLayout: reduce only when literal/column positions align on both sides
(f(id, 2) vs f(4, store_id) is not reducible).
- Gate the deprecated int path on dataType == IntegerType, not the boxed runtime class.
- Fall back to the generalized overload when the deprecated one is absent or returns null.
- Tests: Literal[] reducer end-to-end, deprecated backward-compat, DateType routing,
mismatched argument layout, deprecated-null fallback.
| * position a literal slot aligns with a literal slot (and a non-literal with a non-literal). | ||
| * Literal *values* may differ -- that is what a [[Reducer]] reconciles. | ||
| */ | ||
| private def sameArgumentLayout(other: TransformExpression): Boolean = |
There was a problem hiding this comment.
sameArgumentLayout and isSameFunction (L76-86) are now two separate children.zip(...) walkers. The literal axis differing is deliberate — equality for "same transform" vs. "values may differ, a Reducer reconciles them". But the non-literal handling drifted apart, and sameArgumentLayout ended up strictly weaker:
| position pair | isSameFunction |
sameArgumentLayout |
|---|---|---|
| (literal, literal) | l1 == l2 |
true — intended |
| (transform, transform) | recursive isSameFunction |
case _ => true |
| (col, transform) | isColumnRef && isColumnRef |
case _ => true |
So sameArgumentLayout treats any two non-literals as interchangeable, and its correctness rests entirely on the supportsExpressions gate guaranteeing the single non-literal child is always a plain column ref. On its own it would let two structurally different transforms that share a function name and literal layout reduce (e.g. bucket(4, days(c)) vs bucket(4, hours(c))).
That's fine as the code stands. But since these were a single parameterized helper a commit ago (childrenMatch), it'd be nice to bring sameArgumentLayout back in sync with isSameFunction so the precondition is correct on its own rather than gate-dependent (future-proof):
/** Per-position match; `literalsMatch` is the only axis that varies between callers. */
private def childrenMatch(other: TransformExpression)
(literalsMatch: (Literal, Literal) => Boolean): Boolean =
children.length == other.children.length &&
children.zip(other.children).forall {
case (l1: Literal, l2: Literal) => literalsMatch(l1, l2)
case (t1: TransformExpression, t2: TransformExpression) => t1.isSameFunction(t2)
case (c1, c2) => TransformExpression.isColumnRef(c1) && TransformExpression.isColumnRef(c2)
}
def isSameFunction(other: TransformExpression): Boolean =
function.canonicalName() == other.function.canonicalName() &&
childrenMatch(other)(_ == _)
// reducer precondition: same layout/structure, literal *values* may differ
private def sameArgumentLayout(other: TransformExpression): Boolean =
childrenMatch(other)((_, _) => true)Equally fine to leave as-is if you'd rather not touch it now — mainly flagging so the gate-dependency is a conscious choice.
There was a problem hiding this comment.
@peter-toth I like the idea. Yea, very similar to my parameterized helper created for nested transform. Will make the change.
Thanks.
| * @param otherFunction the other parameterized function | ||
| * @param otherParams literal parameters for the other function | ||
| * @return a reduction function if reducible, null otherwise | ||
| * @since 5.0.0 |
There was a problem hiding this comment.
Version stamp: this should be 4.3.0, not 5.0.0. The companion SPARK-57038 merged to both master (5.0.0) and branch-4.x (4.3.0); assuming this PR is likewise backported, it first ships in 4.3.0, so per the "version of the branch it first ships in" rule both this @since and the @Deprecated(since = "5.0.0") on the deprecated overload below (L118) should read 4.3.0. If this is intended master-only, ignore.
| * @since 5.0.0 | |
| * @since 4.3.0 |
…k; stamp reducer API @SInCE 4.3.0
peter-toth
left a comment
There was a problem hiding this comment.
Thanks @metanil — all my review comments are addressed, LGTM. I plan to merge this Wednesday if there are no further asks or concerns from @szehon-ho or @sunchao.
| // Should only consider column references, not literals. | ||
| val nonLiteralChildren = transform.children.filterNot(_.isInstanceOf[Literal]) | ||
| // We need exactly one column reference per transform. | ||
| nonLiteralChildren.size == 1 && TransformExpression.isColumnRef(nonLiteralChildren.head) |
There was a problem hiding this comment.
[P1] Reducer-compatible transforms must not take the raw partition-key equality fast path.
This widened gate now admits parameterized transforms such as xor(x, 1). With spark.sql.sources.v2.bucketing.pushPartValues.enabled=true, spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled=true, and partially clustered distribution disabled, consider two tables containing x IN (0, 1) partitioned by identity(x) and xor(x, 1). Both scans advertise sorted raw partition keys [0, 1], but partition 0 contains x=0 on the identity side and x=1 on the xor side (partition 1 is reversed). KeyedShuffleSpec.isCompatibleWith treats the reducer-compatible expressions plus equal unreduced keys as directly compatible, so EnsureRequirements skips reducer reconciliation and GroupPartitionsExec; a shuffle-free equality join silently returns zero rows instead of two.
Please restrict the raw-key-equality fast path to identical transform semantics, or force compatible-but-different transforms through reduced-key reconciliation. Add end-to-end SPJ regressions for identity-vs-parameterized and parameterized-vs-parameterized transforms whose raw key arrays coincide but whose physical partitions differ.
|
@metanil this has gone a month without a rebase, and something landed in the meantime that bears directly on it. SPARK-58769 (#58002, merged to master and branch-4.x) dropped the name-based transform comparison, so Could you rebase on master and adjust for it? I will take another pass once it is green. |
|
@peter-toth Thanks for the update. Will update from master adjust the change. I was addressing some naive solution for @sunchao 's P1 feedback, but had created other issues. I'll probably revisiting it a week after next. |
|
@metanil a heads-up on sunchao's P1 about the raw partition-key equality fast path: it is not specific to this PR. The same hole is live on master, through the identity-versus-transform arm SPARK-56182 added. A table partitioned by So there is nothing for you to do here about that P1. Once #58943 lands it falls away on your rebase, and the gate widening in this PR inherits the fix. |
…ntity
Merges upstream/master (1229 commits ahead of this PR's original base) into
the branch. Master independently added its own fix for SPJ co-partitioning
(TransformExpression.reducedWith / hasSameReducedKeys), which supersedes
this PR's original co-partitioning half -- that half is dropped here in
favor of master's version, taken wholesale wherever the merge conflicted
with it.
What master's identity does not do is distinguish a parameterized transform
by anything other than a bucket count. TransformFunctionId is
(canonicalName, numBucketsOpt): populated for bucket, but every other
parameterized transform (e.g. truncate) always reports numBucketsOpt = None.
So truncate(col, 3) and truncate(col, 5) answer `isSameFunction` as the same
transform.
This does not make a truncate reducer dead code -- `reducers()` has no
`isSameFunction` gate at all, so it is consulted regardless of identity.
What the bucket-only identity actually breaks is narrower but still real:
1. `hasSameReducedKeys` is keyed on functionId, so once truncate can reduce,
two unrelated truncate reductions (e.g. truncate(3)/truncate(5), lcm 15,
vs truncate(7)/truncate(11), lcm 77) would collapse onto the same
("truncate", None) identity and compare equal -- two different reduced
key spaces mistaken for one, which drops rows silently. This is latent
only because truncate could not previously reduce at all; the reducer
generalization this PR also carries is what makes it reachable.
2. SPARK-59688 gave `isCompatibleWith` a strict, non-reducing check (asking
`areKeysCompatible` with `allowReduce = false`) that falls through to
`isSameFunction` for a transform-vs-transform pair. Master's bucket-only
identity answers that truncate(3) and truncate(5) are the same transform,
so the strict gate is incomplete for parameterized non-bucket transforms.
Widening identity completes it.
This commit widens identity to close both gaps: TransformFunctionId now
carries the transform's literal parameters (canonicalName,
params: Seq[Literal]) instead of a bucket-specific Option[Int]. This is
strictly more discriminating than master's identity -- every pair master
separated (by bucket count or by name) is still separated, and
parameterized non-bucket transforms are now separated by their own
parameters too.
Kept unchanged from master: reducedWith, reducedKeySpace, hasSameReducedKeys,
reducedTogetherWith, stringArgs, withReference, hasReducedKeys. One gap in
master's identity-vs-transform reducer path is restored rather than dropped:
identityReducer's argsMatchInputTypes gate (this PR's own addition, predating
and orthogonal to reducedWith) is reattached via withReference, since without
it isExpressionCompatible and reducersBothWays could disagree on a
type-mismatched pair -- gate says compatible, reducer refuses -- which keeps
the identity side's raw keys and mis-joins.
All 25 call sites that passed numBucketsOpt positionally (1 production site
in V2ExpressionUtils, which already folds the bucket count into `children`
via NamedTransform; 24 in test files) are updated to fold the literal into
`children` instead. One PR-authored test (nested-transform recursion in
isSameFunction) is rewritten: the widened identity intentionally does not
recurse into nested transforms, since KeyedPartitioning.supportsExpressions
already rejects that shape before identity is ever consulted on it.
Baseline (pre-merge, on the old base): ShuffleSpecSuite 14/14,
KeyGroupedPartitioningSuite + EnsureRequirementsSuite +
ProjectedOrderingAndPartitioningSuite + GroupPartitionsExecSuite 189/189,
4 scalastyle targets clean.
After this commit: ShuffleSpecSuite + TransformExpressionSuite 35/35 (no
bucket assertion edited), the same four sql/core suites 354/354, 4 scalastyle
targets clean. The test that specifically pins this change is
`hasSameReducedKeys discriminates non-bucket (truncate-style) reductions`:
it goes red if functionId is reverted to ignore literal parameters. The
end-to-end truncate SPJ test does not, and should not, depend on this
change: `TransformExpression.reducers` has no `isSameFunction` gate, so it
finds the width reducer regardless of identity -- identity only decides
whether the reduced-key-space bookkeeping (`hasSameReducedKeys`) can tell
two different reductions apart.
…branch-4.x The generalized Literal[]-based reducer overload and the deprecated int overload were stamped @SInCE / @deprecated(since = "4.3.0") when this PR's reducer API work was committed against an older base. branch-4.x has since moved to 4.4.0-SNAPSHOT (confirmed via the GitHub API against apache/spark, since dev/next_version_candidates.py needs a fetchable upstream remote this sandbox's network policy blocks over SSH), so this is a normally-backported SQL change and 4.4.0 is the correct first-release version, not 4.3.0. No other file references the stale "4.3.0" stamp.
…ments Two comments explaining the identity change referred to the removed numBucketsOpt field by name, which a reviewer's `git grep -c numBucketsOpt` would pick up even though the field itself is gone. Reworded to describe the shape instead of naming the removed identifier.
|
@peter-toth Thanks 🙏 . I'll rebase it, and push the necessary change. |
…duce SPARK-59688 (apache#58943) landed on upstream/master since the first merge in this branch's history: it gives areKeysCompatible / isExpressionCompatible an allowReduce parameter, and has isCompatibleWith ask with allowReduce = false, so a pair that still needs a reduce to line up is no longer reported as already lined up. It fixes a live wrong-results bug (identity(id) vs a key-permuting flip_low_bit(id) returned 0 rows of 2 instead of 2). That PR rewrites the same three case arms of isExpressionCompatible this branch's own argsMatchInputTypes gate lives in. Resolved by ANDing both guards (canReduce = allowReduce && canReduceKeys) rather than taking either side wholesale: taking master's side alone would drop the type-safety gate this branch added (see the prior commit), and taking this branch's side alone would drop allowReduce and silently reintroduce the wrong-results bug SPARK-59688 fixes -- with no test on this branch to catch it, since the test that would catch it is SPARK-59688's own, arriving in this same merge. Also fixed two casualties of the same collision: FakeBucket was hoisted from a test-local class to suite scope by the incoming change, so this branch's inputTypes()/produceResult fix (children now carry the bucket count as a Literal, at position 0) had to move with it; and a new SPARK-59688 test constructed FakeBucket-based TransformExpressions with the old Some(numBuckets) positional argument, converted to Literal(numBuckets) in children like every other call site. catalyst/ShuffleSpecSuite + TransformExpressionSuite: 36/36 (was 35, +1 from SPARK-59688's own new isCompatibleWith test). The four sql/core suites: 356/356 (was 354, +2). Four scalastyle targets clean. `allowReduce` appears in all three touched files; `numBucketsOpt` remains at 0.
…not as a bag of literals `isSameFunction` compared the transform's literal children as an unordered bag, so it ignored which argument position each literal occupied and did not look at nested transforms at all. Two consequences: - `truncate(id, 2)` and `truncate(2, store_id)` answered as the same function. `isCompatible` short-circuits on `isSameFunction`, and `sameArgumentLayout` -- the guard for exactly this -- only runs inside `reducer`, so it was unreachable. The pair was read as lined up as it stands and joined with no reduce, pairing partitions that do not hold the same rows. - `bucket(4, years(c))` and `bucket(4, days(c))` answered as the same function, since both literal lists are just [4]. Restore the structural per-position walk: same canonical name, same arity, and per position a literal equals a literal, a nested transform is recursively the same function, and anything else must be a plain column reference on both sides. Column identity stays ignored -- it is reconciled separately, via `keyPositions`. A non-reference slot such as `c + 1` is never "the same", not even against an identical one, since it is not a shape identity can affirm. `TransformFunctionId` is now only the reduced-key-space token stored in `reducedWith` and compared by `hasSameReducedKeys`, which needs a value it can put in a field rather than a predicate. Its one residual coarseness -- a nested transform maps to `None`, so `bucket(4, years(c))` and `bucket(4, days(c))` share a key-space identity -- is documented: it is unreachable for SPJ, since `supportsExpressions` rejects nested shapes, and `isSameFunction` separates them regardless. `KeyGroupedPartitioningSuite` carried a test asserting the bag behaviour, justified by "supportsExpressions rejects this shape before identity would ever be consulted". Identity is also consulted by `isCompatible` and `ValidateRequirements`, so it has to be honest on its own; that test is replaced by one asserting the structural walk.
…ng children Literal parameters (the bucket count, the truncate width) live in `TransformExpression.children` rather than in a field of their own, so every path that rewrites a transform's children has to skip them. `KeyedShuffleSpec.createPartitioning` already did, inline. Alias projection did not. `aliasMap` is built from any `Alias` child, so `SELECT data AS d, 2 AS w` maps `Literal(2) -> w`, and `projectExpression` matches on `aliasMap.contains(e.canonicalized)` for any expression. A `Literal` has no children, so the `containsChild.nonEmpty` fallback does not re-offer the original: the substitution is unconditional. `KeyedPartitioning([truncate(data, 2)])` projected through that becomes `truncate(d, w)`, which - drops the parameter from the transform's identity, - gives the expression two entries in `references`, tripping the single-reference assert in `KeyedShuffleSpec.keyPositions` during planning, and - fails `KeyedPartitioning.supportsExpressions`, so the partitioning is discarded and SPJ is silently lost for any query that projects a matching constant (`0 AS flag`, `1 AS one`). Add `TransformExpression.rewriteColumnSlots` and `columnSlots` so the skip is written once, and route both callers through them: `createPartitioning`, which retargets a transform at the other side's clustering key, and `PartitioningPreservingUnaryExecNode`, which retargets it at an aliased output attribute. A shape with anything other than a single column slot is not projected, matching what `supportsExpressions` admits. This does not arise on master, where the bucket count is held in `numBucketsOpt` outside `children` and `supportsExpressions` requires `children.size == 1`, so no literal is ever a child of a partitioning transform.
…in the reducer gates `inputTypesMatch`, behind `argsMatchInputTypes` and `literalParamsMatchInputTypes`, read each child's `dataType` eagerly. `withReference` retargets a transform at the other side's key, and a `GetStructField` column slot retargeted at a non-struct attribute -- an `identity` column on the other side of a join -- leaves `GetStructField(intCol, 0)`, whose `dataType` raises `INTERNAL_ERROR: GetStructField requires a StructType child`. Both callers are gates whose contract is to answer "not reducible" and let the join shuffle, so the query failed during planning instead. Treat a child whose type cannot be read as not matching. `argsMatchInputTypes`' scaladoc claimed "a gate-admitted child always has a resolvable `dataType`, so the eager read is safe"; that is not true for a retargeted expression, and the doc is corrected.
…; document sorted_bucket With the dedicated `BucketTransform` arm gone from `toCatalystTransformOpt`, bucket columns resolve through `toCatalystOpt`, whose reference arm matched only Spark's `FieldReference`. The removed arm resolved any `NamedReference`, since `resolveRef` only reads `fieldNames`; a connector supplying its own implementation now hit `_LEGACY_ERROR_TEMP_3054`. Match `NamedReference` instead. Also document the one other resolution change: the removed arm was guarded on `sorted.isEmpty`, so a `SortedBucketTransform` with no sort columns used to resolve `bucket` and now resolves its own name, `sorted_bucket`. A connector that does not register that name gets no function and the join shuffles -- correct results, no SPJ. The shape is degenerate and the type `private[sql]`, so it is accepted rather than special-cased.
…one transform identity `isSameFunction` walked the children structurally -- per position, recursing into nested transforms -- while `hasSameReducedKeys` compared `functionId`, which collapsed every non-literal argument to one opaque slot. The two answered "which transforms are the same" differently: `bucket(4, years(c))` and `bucket(4, days(c))` were not the same function, yet shared a reduced key space. That is unreachable today, because `KeyedPartitioning.supportsExpressions` rejects nested transforms, but a later change widening that gate would have inherited a wrong-results hole. Give `TransformFunctionId` the same shape the structural walk compared, as a value: `TransformFunctionId.ArgumentShape` is `Param(literal)`, `Nested(id)` (recursive) or `Column`, and `functionId` is `None` when an argument is none of these (`c + 1`, `cast(c)`), since such a transform is not the same as anything -- not even an identical copy -- and a value is always equal to itself. `isSameFunction` is then `functionId.isDefined && functionId == other.functionId`, a single comparison of a memoized value again, and `hasSameReducedKeys` compares the same identity, so the two cannot drift. `reducedWith.isDefined` is what marks keys as reduced, so a reduce cannot be recorded with a partner that has no identity; `reducedTogetherWith` raises an internal error rather than silently leaving reduced keys looking raw. It is unreachable: its only producer sees partitionings admitted by `supportsExpressions`, which always have one. `hasSameReducedKeys` now requires a defined key space on this side, so two sides without an identity no longer compare equal through `None == None`. Adds tests for nested transforms in `hasSameReducedKeys`, for transforms with no identity, and one asserting that the two checks agree on every pair of a set of shapes.
…make column-slot helpers methods `reducers(other)` only matched both functions as `ReducibleFunction` and forwarded to a private `reducer(thisFunction, thisExpr, otherFunction, otherExpr)`, and `isCompatible` repeated the same match to call it both ways. The four parameters restated what the call site already knew, and let a caller pair one expression's function with the other's parameters. Fold the private method into `reducers(other)`, which now takes the functions from `this` and `other` itself, and write `isCompatible` as what it means: the same function, or reducible in either direction. The dispatch body is unchanged. `rewriteColumnSlots` and `columnSlots` take a `TransformExpression` and rewrite or read its children, the same kind of job as `withReference`, so they move from the companion object to instance methods beside it. The companion keeps the helpers that must accept any `Expression` (`isColumnRef`, `hasReducedKeys`). No behaviour change.
…olumn-argument helper The check that an identity side can be reduced onto a transform side -- `t.withReference(col).argsMatchInputTypes` -- was written four times: twice in `isExpressionCompatible`'s gate and twice in `reducersBothWays`, one per side order, each with a comment that it must agree with the other. They must: `EnsureRequirements` cannot tell a mismatched-type pair from a matched one, so a gate that admitted what the reducer refused would keep raw keys and mis-join. Before the merge with master this was a single function used by both. Restore that as `KeyedShuffleSpec.retargetForIdentity`, which returns the retargeted transform when Spark can evaluate it: the gate admits the pair exactly when it is defined, and the reducer is built from it. `KeyedPartitioning.supportsExpressions` re-implemented `TransformExpression#columnSlots` inline. Call the helper instead, so the gate and alias projection -- which both read "the column argument" -- agree by construction. Trim `projectPartitionExpression`'s scaladoc to point at `rewriteColumnSlots` for why literal parameters must be left alone, keeping only what is specific to alias projection. No behaviour change.
…ssionCompatible The merge resolution hoisted `allowReduce && canReduceKeys` into a local `canReduce`, following an earlier draft of SPARK-59688; the version that landed on master writes it inline at each use. Restore master's form so this PR's diff in `isExpressionCompatible` is only the identity-vs-transform arm. As a side effect `canReduceKeys` no longer reads SQLConf for pairs that never consult it. No behaviour change.
Shorten comments that narrated history, listed every consequence at every site, or repeated a cross-reference at both ends of a shared helper; history is in the commit messages. Restore master's comment on the identity arm of `reducersBothWays`, which this PR only extended. In tests, drop reviewer names and planning labels, and have each comment state what the test guards against rather than what an earlier version did. Comments only.
…reducer The generalized `ReducibleFunction.reducer` Javadoc said array, map, struct and UDT-typed literal parameters are "filtered out by Spark and not passed here", which reads as if the method were called with those parameters removed. Spark instead does not call the method at all for such a pair, and the join shuffles; the same holds for a literal whose type differs from what the function declares. Say so. It also says the two arrays may differ in length, for example a zero-parameter transform reducing onto a one-parameter one. That is true, but Spark only asks when the arguments line up position by position, so `days(ts)` may be offered `truncate(ts, 3)` but never `bucket(4, ts)`. Say that too, and add a test pinning it: `zero_or_one(id)` against `zero_or_one(2, id)` is not reducible, although the fixture's reducer would reduce any pair whose parameter counts differ.
…e same function On master the deprecated `reducer(int, ReducibleFunction, int)` was called only when both sides came from a V2 `bucket` transform, and its Javadoc says so: "This method is for the bucket function", with parameters named `otherBucketFunction` and `otherNumBuckets`. The generalized dispatch widened that to any two functions that each have one int literal and the same argument layout. A connector trusting the documented contract need not check the other function; one that did not could return a gcd reducer for `bucket(4, x)` against a different function taking a single int, and the join would pair partitions that do not hold the same rows. Fall back to the deprecated overload only when both sides are the same function, which restores the guarantee master gave. Bucket against bucket, including Iceberg's, is unchanged. Adds a test with a deprecated-only bucket fixture that does not check the other function: without the check it returns `BucketReducer(2)` for bucket against another literal-first function. Also corrects the comment on the existing legacy-overload test, which still said the deprecated overload is tried first; the generalized one has been tried first since the dispatch was redesigned.
…cated reducer fallback The generalized reducer's Javadoc lists when Spark falls back to the deprecated int overload. Add the condition the previous commit introduced: both sides must be the same function.
|
Thanks @peter-toth for SPARK-59121 and SPARK-59688 — between them they fix the co-partitioning issues this PR originally tried to address (reduced key spaces are now marked and compared via Btw, @peter-toth I modified the your Thanks @sunchao for the original review. The raw-key issue you raised is now in the master with peters fix. |
There was a problem hiding this comment.
Re-checked from scratch through 18f6f8f. Nothing I raised in the earlier rounds regressed, and the one open P1 on the PR, r3679486000, is closed by SPARK-59688 on master: isCompatibleWith asks areKeysCompatible(allowReduce = false), so the identity-vs-transform arm no longer reads a pair whose keys still need reducing as lined up.
The rebase moved a lot, so one line of re-orientation. The co-partitioning half is gone, TransformFunctionId now records each argument's shape in order, and isSameFunction and hasSameReducedKeys are both derived from that one value. That is the right shape: a sidecar field cannot express where a parameter sits, and truncate's width sits at argument 1 while bucket's count sits at argument 0, so moving the parameters into children is what makes the generalization possible at all. Findings 1 and 4 are both "one step too specific" rather than "wrong direction", and I have them as one measured patch.
The findings below are all first raised here.
Blocking
- 1. The literal parameter survives the partitioning projection but not the ordering one (late catch):
dc0afc98b66fixedoutputPartitioningand leftoutputOrderingon the sharedprojectExpression, soProjectExec(id AS pk, 32 AS w)reports the orderingbucket(w#2, pk#1)over a partitioning ofbucket(32, pk#1). Measured, along with the one-clause fix; the two are otherwise out of sync, which is exactly whatGroupPartitionsExec.outputOrdering's coalescing branch filters on. inline:sql/core/src/main/scala/org/apache/spark/sql/execution/AliasAwareOutputExpression.scala:97 - 2. No end-to-end test for the identical-width truncate join (late catch): both new end-to-end truncate tests set
V2_BUCKETING_ALLOW_COMPATIBLE_TRANSFORMS, which is off by default, so they only cover the reducer path. The case the description leads with is the one the default config takes, and it is covered only by a single-tablecheckQueryPlan. inline:sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala:6295
Non-blocking
- 3. A typed-null literal parameter still reaches the generalized overload (late catch): the null clause here keeps it out of the deprecated overload only. Measured:
bucket(null, id).reducers(bucket(4, id))returnsSome(BucketReducer(4)), becauseBucketFunctionreads the null as0andgcd(0, 4)is4. inline:sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala:315 - 4.
isSameFunctionis not reflexive (late catch): identity is re-asking a questionsupportsExpressionsalready owns, so a transform with a non-column-ref argument is not the same function as itself, whileShuffleSpec.isCompatibleWithdocuments reflexivity as assumed. Dropping theOptionfixes it and takes 17 lines with it; measured. inline:sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala:116 - 5. The deprecated-fallback guard compares
name(), notcanonicalName()(new):name()has no uniqueness contract, andcanonicalName()is both the documented one and whatfunctionIdalready uses. The two disagree in both directions. inline:sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala:336 - 6.
IntegerTruncateFunctionshares its canonical name withTruncateFunction(late catch): the two partition data differently, whichBoundFunction.canonicalName's javadoc forbids, sotruncate(strCol, 3)andtruncate(intCol, 3)are the same transform. inline:sql/core/src/test/scala/org/apache/spark/sql/connector/catalog/functions/transformFunctions.scala:494
Alternatives
- 7. Detect an unimplemented overload rather than infer it from a caught
UnsupportedOperationException(late catch): an implemented reducer that throwsUnsupportedOperationExceptionfrom inside is read as unimplemented and silently falls through to the deprecated overload. Optional, with real counter-arguments. inline:sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/TransformExpression.scala:322
Minor
- 8. Four comments still describe a dispatch order and a guard shape that changed (new): one of them contradicts the assertion in the test that uses its fixture. inline:
sql/core/src/test/scala/org/apache/spark/sql/connector/catalog/functions/transformFunctions.scala:670 - 9. The two
ReducibleFunctiondefaults can recurse into each other (late catch): the natural migration shim delegates back, and the resultingStackOverflowErroris fatal, soTrydoes not catch it. inline:sql/catalyst/src/main/java/org/apache/spark/sql/connector/catalog/functions/ReducibleFunction.java:158 - 10. Two gaps in the description (new): the fallback rule omits the condition
22cd902c0d4added, that both sides must be the same function. And "Paths that rewrite a transform's children (KeyedShuffleSpec.createPartitioning, alias projection) go through one helper" does not hold for the ordering alias projection, see finding 1.
| * only the column argument is projected: `aliasMap` also maps aliased literals, so `2 AS w` would | ||
| * otherwise turn `truncate(data, 2)` into `truncate(d, w)`. | ||
| */ | ||
| private def projectPartitionExpression(expr: Expression): LazyList[Expression] = expr match { |
There was a problem hiding this comment.
Finding 1. dc0afc98b66 closed this for outputPartitioning and left outputOrdering on the shared projectExpression. DataSourceV2ScanExecBase.outputOrdering derives the ordering from the same k.expressions (v2BucketingPartitionKeyOrderingEnabled is on by default), and AliasAwareQueryOutputOrdering.outputOrdering projects it through projectExpression directly, so the literal is still substituted there.
Measured on this head, with a dummy leaf reporting KP([bucket(32, id)]) and outputOrdering = [bucket(32, id) ASC], under ProjectExec(id AS pk, 32 AS w):
partitioning: keyedpartitioning(transformexpression(bucket, 32, pk#1), ...)
ordering: List(transformexpression(bucket, w#2, pk#1) ASC NULLS FIRST)
ordering refs: {w#2, pk#1}
ExpressionSet(partitioning.expressions).contains(ordering.head.child): false
Control, same plan without 32 AS w: ordering transformexpression(bucket, 32, pk#3), and contains is true.
That last line is the consequence. GroupPartitionsExec.outputOrdering's coalescing branch keeps exactly the sort orders whose expression is a partition key expression:
val keyExprs = ExpressionSet(keyedPartitionings.flatMap(_.expressions))
child.outputOrdering.filter(order => keyExprs.contains(order.child))and its comment states the premise this breaks: "The child's outputOrdering should already be in sync with the partitioning". With the literal substituted the filter drops the key-derived ordering, so a downstream sort that would have been eliminated is kept. A GroupPartitionsExec over a ProjectExec over a BatchScanExec is the ordinary shape of a join leg, so this is the default-config path. The claim itself stays value-correct, since w is the constant 32 in every row, which is why nothing fails. It just stops matching.
The fix I would suggest is central rather than per-site, at sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/AliasAwareOutputExpression.scala:74:
case e: Expression if !e.isInstanceOf[Literal] && aliasMap.contains(e.canonicalized) =>
aliasMap(e.canonicalized).toSeq ++ (if (e.containsChild.nonEmpty) Seq(e) else Seq.empty)A Literal then stays put for every consumer, which is what you want everywhere: an alias over a constant is not another name for the same value, it is a column that happens to hold it. projectPartitionExpression can then go away entirely, 18 lines lighter.
I ran it, together with finding 4's change, on this head. Green: TransformExpressionSuite, ShuffleSpecSuite, KeyGroupedPartitioningSuite, ProjectedOrderingAndPartitioningSuite, GroupPartitionsExecSuite, EnsureRequirementsSuite, PlannerSuite (479 tests), plus ReplaceIntegerLiteralsWithOrdinalsDataframeSuite and RemoveRedundantSortsSuite. No test in sql/core partitions or orders by a bare literal, so nothing depended on the old behaviour.
Either way the new ProjectedOrderingAndPartitioningSuite test wants the ordering half asserted next to the partitioning one, since that is the pair that has to stay in sync:
val orderedChild = DummyLeafExecWithPartitioningAndOrdering(
output = Seq(id),
partitioning = KeyedPartitioning(Seq(bucketExpr), keys1d),
ordering = Seq(SortOrder(bucketExpr, Ascending)))
val orderedProject = ProjectExec(Seq(pk, w), orderedChild)
val projectedOrdering = orderedProject.outputOrdering
assert(projectedOrdering.head.child.asInstanceOf[TransformExpression].children.head
=== Literal(32),
"the literal parameter must survive the ordering projection too")
val projectedKp = orderedProject.outputPartitioning.asInstanceOf[KeyedPartitioning]
assert(ExpressionSet(projectedKp.expressions).contains(projectedOrdering.head.child),
"the projected ordering must stay in sync with the projected partitioning")If you prefer to keep the change narrow instead, give AliasAwareQueryOutputOrdering the same TransformExpression case projectPartitionExpression has.
| assert(shuffles.isEmpty, | ||
| "truncate(3) vs truncate(5) should avoid shuffle via the width reducer, " + | ||
| "but a shuffle was planned") | ||
| checkAnswer(df, Seq(Row(0, 10), Row(1, 20), Row(2, 30))) |
There was a problem hiding this comment.
Finding 2. Both new end-to-end truncate tests set V2_BUCKETING_ALLOW_COMPATIBLE_TRANSFORMS, which is off by default, so between them they only cover the reducer path. The case the description leads with, identical widths matched by isSameFunction with no reducer at all, is covered only by the single-table checkQueryPlan at line 549. That is the path the default config takes and the one most users will hit.
It works. The test is one copy of this one without the withSQLConf. Measured on this head: 0 shuffles and the right rows, with allowCompatibleTransforms left at false.
test("SPARK-50593: identical truncate widths trigger SPJ under the default config") {
val table1 = "trunc_same_a"
val table2 = "trunc_same_b"
val partitions = Array(
Expressions.apply("truncate", Expressions.column("data"), Expressions.literal(3)))
createTable(table1, columns, partitions)
sql(s"INSERT INTO testcat.ns.$table1 VALUES " +
"(0, 'apple', CAST('2022-01-01' AS timestamp)), " +
"(1, 'grape', CAST('2021-01-01' AS timestamp)), " +
"(2, 'orange', CAST('2020-01-01' AS timestamp))")
createTable(table2, columns, partitions)
sql(s"INSERT INTO testcat.ns.$table2 VALUES " +
"(10, 'apple', CAST('2022-01-01' AS timestamp)), " +
"(20, 'grape', CAST('2021-01-01' AS timestamp)), " +
"(30, 'orange', CAST('2020-01-01' AS timestamp))")
// No withSQLConf: identical widths are the same transform, so this has to work on the
// default config, where allowCompatibleTransforms is off.
val df = sql(
s"""
|${selectWithMergeJoinHint(table1, table2)}
|$table1.id AS left_id, $table2.id AS right_id
|FROM testcat.ns.$table1 JOIN testcat.ns.$table2
|ON $table1.data = $table2.data
|ORDER BY $table1.id
|""".stripMargin)
assert(collectShuffles(df.queryExecution.executedPlan).isEmpty,
"identical truncate widths are the same function, so no shuffle is needed")
checkAnswer(df, Seq(Row(0, 10), Row(1, 20), Row(2, 30)))
}It cannot pass on base: supportsExpressions required children.size == 1 there, which truncate(data, 3) fails, and the pre-PR version of the test at line 549 asserted UnknownPartitioning(0) for exactly that reason.
| // Only a single non-null IntegerType parameter per side may use the deprecated int overload; a | ||
| // typed null would otherwise be read as 0. | ||
| def isSingleInt(p: Array[V2Literal[_]]): Boolean = { | ||
| p.length == 1 && p(0).dataType == IntegerType && p(0).value() != null |
There was a problem hiding this comment.
Finding 3. This clause keeps a typed null out of the deprecated overload, which is what r3503138262 asked for. The generalized overload still receives it: literalParamsMatchInputTypes compares only dataType, and noComplexLiteralParams looks only at dataType too, so Literal(null, IntegerType) passes both gates.
All three bundled reducers then read the value with .value().asInstanceOf[Int], and in Scala null.asInstanceOf[Int] is 0. Measured on this head:
bucket(null, id).isSameFunction(bucket(4, id)): false
bucket(null, id).literalParamsMatchInputTypes: true
bucket(null, id).reducers(bucket(4, id)): Some(BucketReducer(4))
bucket(null, id).isCompatible(bucket(4, id)): true
BucketFunction.reducer read the null as 0, computed gcd(0, 4) = 4, and returned a reducer, so Spark declares the pair reducible and reduces the null side by % 4. That is the same fabricated zero as the deprecated path, one overload up.
A null parameter belongs with the other things the gate refuses, rather than with something the javadoc warns about. Then isSingleInt needs no null clause of its own:
if (!sameArgumentLayout(other) ||
!literalParamsMatchInputTypes || !other.literalParamsMatchInputTypes ||
!noComplexLiteralParams || !other.noComplexLiteralParams ||
literalChildren.exists(_.value == null) ||
other.literalChildren.exists(_.value == null)) {
return None
}If you would rather let a null through, the javadoc needs to say so, because right now it tells an implementor to read value() as "Spark's internal representation" and only warns about array lengths.
| val shapes = children.map { | ||
| case l: Literal => Some(ArgumentShape.Param(l)) | ||
| case t: TransformExpression => t.functionId.map(ArgumentShape.Nested(_)) | ||
| case c if TransformExpression.isColumnRef(c) => Some(ArgumentShape.Column) |
There was a problem hiding this comment.
Finding 4. This arm gives up the identity when an argument is not a literal, a nested transform or a plain column ref, so functionId is None and isSameFunction answers false for such a transform against itself. ShuffleSpec.isCompatibleWith documents the opposite at sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/physical/partitioning.scala:1586: "Note that Spark assumes this to be reflexive, symmetric and transitive". The new TransformExpressionSuite test pins the non-reflexivity ("never the same, not even as itself").
It is reachable under v2BucketingShuffleEnabled. KeyedShuffleSpec.createPartitioning builds the shuffled side's expression from the other side's clustering key, so a join ON t1.id = t2.b + 1 over a bucket(4, id) table yields bucket(4, b + 1) on the shuffled side, paired with bucket(4, id) on the keyed one. keyPositions reconciles them, one reference each, and the two are co-located by construction, since createPartitioning built the expression for exactly that. Master's isSameFunction said true for that pair; this head says false.
The root cause is that identity is being asked two questions. "Which transform is this" is identity's; "may SPJ reason about these argument shapes" is KeyedPartitioning.supportsExpressions', which already asks it and is the load-bearing gate on the scan. Once identity stops re-asking it, the whole thing gets shorter:
lazy val functionId: TransformFunctionId = {
import TransformFunctionId.ArgumentShape
TransformFunctionId(function.canonicalName(), children.map {
case l: Literal => ArgumentShape.Param(l)
case t: TransformExpression => ArgumentShape.Nested(t.functionId)
case _ => ArgumentShape.Slot
})
}with Column renamed to Slot, since it now covers any non-parameter argument. What that discriminates is unchanged where it matters:
truncate(id, 2)andtruncate(2, store_id)stay apart, theParampositions differbucket(4, years(c))andbucket(4, days(c))stay apart,Nestedis still matched abovebucket(32, c)andbucket(16, c)stay apart
and the knock-on simplifications fall out:
isSameFunctionbecomesfunctionId == other.functionIdreducedKeySpacegoes back from the for-comprehension toreducedWith.map(partner => Set(functionId, partner))reducedTogetherWithloses itsSparkException.internalErrorarm, which its own comment already called unreachable, so it is a one-liner again- the
SparkExceptionimport goes
Net -17 lines in TransformExpression.scala.
On whether the wider Slot is safe. Identity is column-blind by construction, and the safety comes from keyPositions plus supportsExpressions, not from identity. For two non-column-slot transforms to be compared at all, both have to be createPartitioning products, since the scan gate refuses such a shape, and those are built from one spec onto one key set, so "same function" is the right answer there. I could not construct a case where the wider slot co-locates rows that do not belong together; if you see one, that settles it the other way.
I ran this on this head. Green: TransformExpressionSuite, ShuffleSpecSuite (41), then KeyGroupedPartitioningSuite, ProjectedOrderingAndPartitioningSuite, GroupPartitionsExecSuite, EnsureRequirementsSuite, PlannerSuite (438). Four assertions had to flip, and three of them are the ones that enshrine the non-reflexivity:
TransformExpressionSuite, "a transform with a non-reference argument has no identity" -- becomes "a non-parameter argument is a slot, whatever its shape"- the
plusOneassertion in "transform identity is per-argument-position" - the
addassertion inKeyGroupedPartitioningSuite's "isSameFunction is a structural per-position walk" - the
shapes.forall(_.functionId.isDefined)fixture guard in "isSameFunction and hasSameReducedKeys agree on every pair", which can instead gain a non-reference shape to the fixture list, since the agreement now holds over those too
Pushed both changes on top of 18f6f8f in case cherry-picking is easier than redoing them: SPARK-50593-simplify-identity-and-alias-projection on my fork. Two independent commits, so you can take either on its own:
4c51b96-- this finding, the totalfunctionIdplus the four test assertions that flipe93ae9b-- finding 1, the one-clauseprojectExpressionrule,projectPartitionExpressiondeleted, and the ordering assertion added toProjectedOrderingAndPartitioningSuite
75 insertions and 78 deletions over 6 files. Same content I ran the suites against; take it, adapt it or ignore it.
There was a problem hiding this comment.
Please hold off on this one, and on 4c51b96, for now. I found a case where the total functionId looks wrong, and I'm still looking into it. I'll follow up here.
There was a problem hiding this comment.
Following up on the hold above. I was wrong here: a transform over b + 1 must never be the same as one over a column, and your None is what guarantees that on this head.
The total functionId makes bucket(4, b + 1) the same as bucket(4, x), and that loses rows:
KeyedShuffleSpec.createPartitioningshuffles a side ontobucket(4, b + 1)when its join key isb + 1.- A later join that clusters that side on the bare
bpairs it with abucket(4, x)table, becausekeyPositionsmaps a partition expression to a cluster key through its reference.
With 4c51b96 on top of this head, t1 JOIN plain p ON t1.id = p.b + 1 JOIN t3 ON p.b = t3.x returns 0 of 7 rows. This head returns all 7.
Master has the same bug, since its isSameFunction ignores the arguments. I opened SPARK-59887 and #59165 for it. It makes isSameFunction and reducers refuse a transform whose argument is not a bare Attribute, keeps such a layout from being the one other children are shuffled onto, and adds tests for the shapes it fixes. Could you rebase onto it once it lands? Three things to watch in the conflict:
- The new
argumentsAreAttributesskips literal children, so it carries over to your model as it is. - It also refuses a struct field argument, which your
Optiontreats as aColumn. A side shuffled ontobucket(4, s.a)pairs with one onbucket(4, s.b)in a join ons, and loses rows the same way. Please keep that refusal, either through the new guard or by makingArgumentShape.Columna bareAttributeonly. - The first new test also asserts two shuffles. This head gives the right rows but plans three, since it has no
canCreatePartitioningclause yet, so it passes only once that clause is in.
Please don't take 4c51b96. e93ae9b (finding 1) is unrelated and still stands.
There was a problem hiding this comment.
Thanks @peter-toth for the update. I was exploring both findings (1 and 4) yesterday. I'll for wait for your PR #59165 and rebase it here. I'll try to preemptively pull it before the merge to see how it fits here, and be ready.
| // The deprecated overload is documented for bucket against bucket, so only offer it a | ||
| // pair of the same function. | ||
| case Unimplemented if isSingleInt(thisParams) && isSingleInt(otherParams) && | ||
| function.name() == other.function.name() => |
There was a problem hiding this comment.
Finding 5. name() has no uniqueness contract. Function.name() is a display name, and Spark's own fixtures return it both bare (BucketFunction, "bucket") and catalog-qualified (ShuffleSpecSuite.FakeBucket, "test.fakeBucket"). canonicalName() is the one documented for this question, and it is what functionId already compares two hundred lines up.
The two disagree in both directions.
- Over-admits. Two different bound functions may share a
name(). This PR's ownUnboundTruncateFunctionbindsTruncateFunctionorIntegerTruncateFunctionby column type, and both return"truncate", so a pair of those passes this clause. - Under-admits.
BoundFunction.canonicalName's javadoc exists for the case where one implementation is loaded under two catalog names: "loaded using different names, liketest.func_nameandprod.func_name". A connector whosename()carries the catalog then fails this clause for two equivalent functions, and loses the fallback the commit was protecting.
| function.name() == other.function.name() => | |
| function.canonicalName() == other.function.canonicalName() => |
| override def inputTypes(): Array[DataType] = Array(IntegerType, IntegerType) | ||
| override def resultType(): DataType = IntegerType | ||
| override def name(): String = "truncate" | ||
| override def canonicalName(): String = name() |
There was a problem hiding this comment.
Finding 6. TruncateFunction above returns the same "truncate" here, and the two partition data differently: one takes a string prefix, the other snaps an int down to a grid. BoundFunction.canonicalName's javadoc rules that out: "Two functions that partition data differently must not return the same name; Spark may otherwise treat unrelated data as co-partitioned."
isSameFunction keys on the canonical name plus the argument shapes, so truncate(strCol, 3) and truncate(intCol, 3) are the same transform today. Nothing in the suite joins those two tables, so it is latent. But these fixtures are what a connector author reads as the reference implementation of the API this PR is generalizing, and the contract they break is the one SPJ rests on.
| override def canonicalName(): String = name() | |
| override def canonicalName(): String = "truncate(int)" |
and "truncate(string)" on TruncateFunction. Leave both name()s as "truncate", so the function catalog still resolves them the same way.
| def probe(call: => Reducer[_, _]): Outcome = Try(Option(call)) match { | ||
| case Success(Some(r)) => Reducible(r) | ||
| case Success(None) => NotReducible | ||
| case Failure(_: UnsupportedOperationException) => Unimplemented |
There was a problem hiding this comment.
Finding 7. Optional, and a design point rather than a defect.
The dispatch reads "the connector did not implement this overload" off a caught UnsupportedOperationException, and the new javadoc now states that convention. But it is not a property of the overload, it is a property of one call. An implemented generalized reducer that throws UnsupportedOperationException from inside, say a connector operation it cannot serve for these parameters or a require in a helper, is read as unimplemented. The pair then falls through to the deprecated overload, which may hand back a reducer for parameters the generalized one refused. Silent, and the "implements no reducer" warning that would otherwise appear is suppressed by the very fallback that fired.
Whether the method is overridden is answerable directly, and there is precedent right next door: V2ExpressionUtils.findMethod already does this for ScalarFunction's magic method. Sketch:
// Memoized per BoundFunction class, since this is asked per pair.
private def implementsGeneralizedReducer(f: ReducibleFunction[_, _]): Boolean =
f.getClass
.getMethod("reducer", classOf[Array[V2Literal[_]]], classOf[ReducibleFunction[_, _]],
classOf[Array[V2Literal[_]]])
.getDeclaringClass != classOf[ReducibleFunction[_, _]]Unimplemented would then be decided before the call rather than inferred from it, and an UnsupportedOperationException out of an implemented reducer would join Threw.
Counter-arguments, so this is a choice and not a request. Reflection on a connector class costs something, even memoized per class. A Scala object implementing the Java interface makes getDeclaringClass less obvious than it reads. And the UOE convention is what the other default methods in this package already use, so changing it only here makes the package inconsistent.
Fine to leave as is. If you do, the javadoc should say that an UnsupportedOperationException thrown from inside an implemented reducer is indistinguishable from not implementing it, since an implementor cannot otherwise guess that.
There was a problem hiding this comment.
Yes, this was the known issue, rare but not impossible.
I didn't want to use reflection for it, but as you pointed out, there's already a history of using it, I might rethink this.
| /** | ||
| * A bucket-like function implementing BOTH reducer overloads: the deprecated int overload succeeds | ||
| * (returns a reducer), while the generalized Literal[] overload throws. Used to verify the | ||
| * single-int dispatch is lazy -- it must not invoke (and log the throw from) the generalized |
There was a problem hiding this comment.
Finding 8. This one contradicts the test that uses it. The fixture doc says the dispatch "must not invoke (and log the throw from) the generalized overload once the deprecated one already produced a reducer", while the test "a throwing generalized overload is surfaced, not masked by the deprecated one" asserts the reverse: the generalized throw is authoritative and there is no fallback.
Three more sites still describe a shape that changed.
KeyGroupedPartitioningSuite.scala:6343says "the single-int dispatch tries reducer(int, ...) first (UOE), then the Literal[]". It has been generalized-first since the dispatch was redesigned, andBucketFunctionoverrides only the new overload, so no fallback is involved in that test at all.transformFunctions.scala:313,DualApiBucketFunction, says it verifies "that the dispatch falls back to the generalized overload when the deprecated one returns null". The deprecated overload is never consulted there, so its null is irrelevant, which is what the test body now says.transformFunctions.scala:576and:595,StructBackedUDTandUdtParamFunction, describenoComplexLiteralParamsas value-based and call the UDT "the case aDataType-based container check misses but a value-based one catches". The implemented guard isDataType-based and lists_: UserDefinedType[_]explicitly, so it catches it head-on.
22cd902c0d4 fixed a fourth site of this kind. These are the rest.
| */ | ||
| default Reducer<I, O> reducer(ReducibleFunction<?, ?> otherFunction) { | ||
| throw new UnsupportedOperationException(); | ||
| return reducer(new Literal<?>[0], otherFunction, new Literal<?>[0]); |
There was a problem hiding this comment.
Finding 9. With this delegation the two defaults can call each other. A connector migrating to the generalized overload writes the natural shim for its zero-parameter transforms:
@Override
public Reducer<I, O> reducer(
Literal<?>[] thisParams, ReducibleFunction<?, ?> other, Literal<?>[] otherParams) {
if (thisParams.length == 0 && otherParams.length == 0) {
return reducer(other); // keep the old no-parameter logic
}
...
}reducer(other) is not overridden, so it lands back here, which calls the override again. TransformExpression.reducers calls reducer(otherFunction) for exactly the zero-parameter pair, so the recursion is reachable, and the resulting StackOverflowError is a VirtualMachineError: Try's NonFatal does not catch it, so it takes the driver thread down rather than costing a shuffle.
An @implNote is enough, since the direction is not guessable from the signatures:
* @implNote This default delegates to
* {@link #reducer(Literal[], ReducibleFunction, Literal[])}, so an implementation of that
* overload must not call back into this method.
There was a problem hiding this comment.
Ya, I wanted to avoid adding implementation detail in javadoc, but good idea to add was part of @implNote. Thanks.
There was a problem hiding this comment.
@peter-toth Just a note. Looks like Spark doesn't have @implNote build support. I'll add that one too in the build.
Let me know if we don't want to change build file in this PR.
There was a problem hiding this comment.
Thanks for flagging it. That's on me: I suggested the tag without checking that the build doesn't register it. I'd rather not introduce a new Javadoc tag for a single note, so please drop the build changes and keep it as a plain paragraph instead:
* <p>
* This default delegates to
* {@link #reducer(Literal[], ReducibleFunction, Literal[])}, so an implementation of that
* overload must not call back into this method.
`projectExpression` replaced any expression found in `aliasMap`, literals included, and a `Literal` has no children, so the original was not offered alongside the alias. With transform parameters in `children`, `SELECT data AS d, 2 AS w` over a `truncate(data, 2)` partitioning or ordering produced `truncate(d, w)`. The partitioning path worked around this for top-level transforms only; the ordering path, which a scan derives from its partition expressions by default, did not, so after such a projection the two disagreed and a matching ordering was dropped. Leave literals alone in `projectExpression` itself: a literal means the same in any scope, so replacing it with an alias adds nothing. This covers partitionings, orderings and nested transforms in one place, and the partitioning-only workaround is removed. Adds tests for the ordering case and for a nested transform.
The case stands for a column reference -- an `Attribute` or a `GetStructField` chain on one, as decided by `TransformExpression.isColumnRef` -- not a column value. Name it accordingly. No behaviour change.
…on identity and test coverage - Refuse a pair with a null literal parameter before asking either reducer overload. The null was kept out of the deprecated int overload only, so the generalized one could still read it as 0: `bucket(null, id)` against `bucket(4, id)` returned `BucketReducer(4)`. `isSingleInt` no longer needs its own null clause, and the Javadoc lists a null literal among the cases Spark does not ask. - Compare `canonicalName()` rather than `name()` in the deprecated-fallback guard. `name()` is a display name with no uniqueness contract; `canonicalName()` is the documented identity and what `functionId` compares. - Give the two truncate test fixtures distinct canonical names, `truncate(string)` and `truncate(int)`: they partition data differently, which `BoundFunction.canonicalName` forbids for a shared name. Both keep `name()` "truncate", so catalog resolution is unchanged. - Add an end-to-end test for identical truncate widths on the default config, where `allowCompatibleTransforms` is off; the existing truncate tests only covered the reducer path. Each fix comes with a test that fails without it.
… fixture comments The no-parameter `ReducibleFunction.reducer(otherFunction)` default delegates to the generalized overload, so an implementation of the generalized overload that calls back into it recurses until `StackOverflowError`, which is fatal and not caught by the dispatch. Document it with `@implNote`. Spark's Javadoc runs with `-Xdoclint:all`, under which an unregistered `@implNote` is an error, so register the tag alongside `example` and `note` in both the sbt unidoc options and the Maven javadoc plugin. Also correct four test comments that still described the dispatch before it became generalized-first, or the complex-literal guard as value-based: one contradicted the assertion of the test that uses its fixture.
|
Thanks @peter-toth. I've addressed all of your findings except 4 (as per you instruction) and 7. |
|
@metanil, #59165 is merged (master
|
What changes were proposed in this pull request?
This PR adds Storage Partitioned Join (SPJ) support for the
truncatepartition transform. The approach generalizes theReducibleFunctionAPI to accept arbitrary transform parameters, so SPJ can reason about any parameterized transform (bucket, truncate, future ones) through one code path.Key changes:
ReducibleFunction.reducer(Literal<?>[], ReducibleFunction, Literal<?>[]). Parameters are V2Literals, so a connector sees each value together with itsDataType. Spark tries this overload first and falls back to the deprecatedreducer(int, ..., int)only when it is not implemented, both sides are the same function, and each side has a single integer parameter, so existing connectors (e.g., Iceberg, which implements only the old API) continue to work unchanged.TransformExpressionrefactor: literal parameters (e.g., bucketnumBuckets, truncatewidth) now live insidechildrenrather than a bespokenumBucketsOpt: Option[Int]field, andV2ExpressionUtilsresolves bucket through the same generic path as every other transform.TransformFunctionIdrecords each argument's shape in order (literal, nested transform, or column), and bothisSameFunctionandhasSameReducedKeyscompare it, sotruncate(col, 3)andtruncate(col, 5)are different transforms and reduce onto different key spaces.KeyedShuffleSpec.createPartitioningretargets a transform through a helper that leaves literal parameters untouched, and alias projection never substitutes a literal, for partitionings and orderings alike. So e.g.SELECT data AS d, 2 AS wover atruncate(data, 2)partitioning keeps SPJ, and its ordering still matches the partitioning.Why are the changes needed?
Today a join on tables partitioned by
truncate(col, N)always shuffles, even when both sides share identical partitioning. The write-side was fixed by SPARK-40295 (Allow v2 functions with literal args in write distribution and ordering), but the read/join side was never enabled.Previous work in #49211 (@szehon-ho) explored direct support for transforms with literal arguments by adjusting the SPJ paths to recognize them. This PR generalizes the reducer API so the compatibility check is function-agnostic.
It builds on two changes now on master. SPARK-59121 records when a join has reduced both sides onto a common key space (
TransformExpression.reducedWith) and compares those key spaces by transform identity; SPARK-59688 (#58943) stops the join pairing raw keys that still need reducing, the issue raised earlier in this review. Both compare transforms by identity, which is why identity has to carry the parameters once non-bucket transforms can reduce.Does this PR introduce any user-facing change?
Yes, for connector/catalog authors:
ReducibleFunction.reducer(Literal<?>[], ReducibleFunction, Literal<?>[])(@since 4.4.0);reducer(int, ReducibleFunction, int)is deprecated.LegacyBucketFunctiontest fixture.For end users, queries joining tables partitioned by compatible
truncatetransforms (identical widths, or reducible pairs liketruncate(3)andtruncate(5)) now avoid shuffle via SPJ.How was this patch tested?
43 new tests across
KeyGroupedPartitioningSuite,TransformExpressionSuiteandProjectedOrderingAndPartitioningSuite: truncate SPJ end to end, transform identity (parameter values and positions, nested transforms), parameters surviving alias projection, and the reducer dispatch and backward-compatibility paths.cc @szehon-ho @aokolnychyi @sunchao @peter-toth
Was this patch authored or co-authored using generative AI tooling?
Yes.
Generated-by: Claude Code (Opus 5)