Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
22 commits
Select commit Hold shift + click to select a range
011f5a5
feat: [SPARK-59277][SQL][HIVE] Support first-class CHAR/VARCHAR in Hi…
srielau Sep 7, 2026
97ebb3f
fix: [SPARK-59277][SQL][HIVE] Enforce CHAR/VARCHAR boundaries
srielau Sep 8, 2026
3b611b3
fix: [SPARK-59277] preserve analyzed Hive conversion types
srielau Sep 9, 2026
2294f32
test: [SPARK-59277] update Hive reflection expectation
srielau Sep 10, 2026
c770943
fix: [SPARK-59277] preserve nested Hive conversion types
srielau Sep 11, 2026
3b0a504
fix: [SPARK-59277] cover Hive CHAR conversion review gaps
srielau Sep 16, 2026
8961995
test: [SPARK-59277] cover Hive inspector edge cases
srielau Sep 16, 2026
e239606
refactor: [SPARK-59277] simplify Hive CHAR conversion
srielau Sep 16, 2026
484bd6b
fix: [SPARK-59277] snapshot HiveGenericUDF CHAR result types
srielau Sep 16, 2026
21df1f6
fix: [SPARK-59277] snapshot Hive types on the product and fail TRANSF…
srielau Sep 16, 2026
49396dd
fix: [SPARK-59277] limit SerDe rewrite and check Hive runtime types
srielau Sep 17, 2026
9ceec6e
fix: [SPARK-59277] enforce nested TRANSFORM CHAR semantics
srielau Sep 17, 2026
c825068
fix: [SPARK-59277] physicalize nested TRANSFORM CHAR without SerDe
srielau Sep 17, 2026
808127c
fix: [SPARK-59277] rewrite nested JSON map keys recursively
srielau Sep 17, 2026
0cc7cdb
fix: [SPARK-59277] address Hive review feedback
srielau Sep 19, 2026
05da5e1
fix: [SPARK-59277] fresh map builder per row in no-SerDe key restoration
srielau Sep 19, 2026
d904402
fix: [SPARK-59277] tighten restorer API and force single-partition test
srielau Sep 19, 2026
f1230d5
refactor: [SPARK-59277] move script TRANSFORM CHAR/VARCHAR support to…
srielau Sep 21, 2026
a0fb439
fix: [SPARK-59277] validate Hive map keys after CHAR conversion and t…
srielau Sep 21, 2026
cec2eff
fix: [SPARK-59277] check struct field names, remove orphaned in-place…
srielau Sep 21, 2026
ae383a1
fix: [SPARK-59277] test CHAR map-key collision with raw Java inspecto…
srielau Sep 21, 2026
bf11663
fix: [SPARK-59277] cover writable inspector and assert LAST_WIN value…
srielau Sep 21, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
268 changes: 231 additions & 37 deletions sql/hive/src/main/scala/org/apache/spark/sql/hive/HiveInspectors.scala

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -111,34 +111,36 @@ class HiveSimpleUDFEvaluator(
}
}

class HiveGenericUDFEvaluator(
funcWrapper: HiveFunctionWrapper, children: Seq[Expression])
extends HiveUDFEvaluatorBase[GenericUDF](funcWrapper, children) {

// SPARK-58792: copied expression nodes (e.g. via withNewChildrenInternal) share one
// HiveFunctionWrapper, whose cached GenericUDF instance is mutable: initialize()
// rewrites its converters and output holders based on the arguments of whichever
// copy initialized it last. Give every evaluator its own clone so copied nodes
// cannot corrupt each other.
@transient
override lazy val function: GenericUDF =
HiveFunctionRegistryUtils.cloneGenericUDF(funcWrapper.createFunction[GenericUDF]())

@transient
private lazy val argumentInspectors = children.map(toInspector).toArray
private[hive] object HiveGenericUDFEvaluator extends HiveInspectors {

/**
* Driver-side Hive initialization for `SELECT hive_udf(...)`. Returns the Catalyst type
* (for example CHAR(5) from a CHAR inspector). `HiveGenericUDF.apply` stores it.
*/
def inferReturnType(
funcWrapper: HiveFunctionWrapper,
children: Seq[Expression]): DataType = {
val function =
HiveFunctionRegistryUtils.cloneGenericUDF(funcWrapper.createFunction[GenericUDF]())
inspectorToDataType(initialize(function, children.map(toInspector).toArray))
}

@transient
lazy val returnInspector = {
def initialize(
function: GenericUDF,
argumentInspectors: Array[ObjectInspector]): ObjectInspector = {
// Inline o.a.h.hive.ql.udf.generic.GenericUDF#initializeAndFoldConstants, but
// eliminate calls o.a.h.hive.ql.exec.FunctionRegistry to avoid initializing Hive
// built-in UDFs.
val oi = function.initialize(argumentInspectors)
val udfType = function.getClass.getAnnotation(classOf[HiveUDFType])
val isDeterministic =
udfType != null && udfType.deterministic() && !udfType.stateful()
// If the UDF depends on any external resources, we can't fold because the
// resources may not be available at compile time.
if (function.getRequiredFiles == null && function.getRequiredJars == null &&
argumentInspectors.forall(ObjectInspectorUtils.isConstantObjectInspector) &&
!ObjectInspectorUtils.isConstantObjectInspector(oi) &&
isUDFDeterministic &&
isDeterministic &&
ObjectInspectorUtils.supportsConstantObjectInspector(oi)) {
val argumentValues: Array[DeferredObject] = argumentInspectors.map { argumentInspector =>
new GenericUDF.DeferredJavaObject(
Expand All @@ -155,16 +157,42 @@ class HiveGenericUDFEvaluator(
oi
}
}
}

private[hive] class HiveGenericUDFEvaluator(
funcWrapper: HiveFunctionWrapper,
children: Seq[Expression],
catalystReturnType: DataType)
extends HiveUDFEvaluatorBase[GenericUDF](funcWrapper, children) {

// SPARK-58792: copied expression nodes (e.g. via withNewChildrenInternal) share one
// HiveFunctionWrapper, whose cached GenericUDF instance is mutable: initialize()
// rewrites its converters and output holders based on the arguments of whichever
// copy initialized it last. Give every evaluator its own clone so copied nodes
// cannot corrupt each other.
@transient
override lazy val function: GenericUDF =
HiveFunctionRegistryUtils.cloneGenericUDF(funcWrapper.createFunction[GenericUDF]())

@transient
private lazy val argumentInspectors = children.map(toInspector).toArray

@transient
lazy val returnInspector = {
val inspector = HiveGenericUDFEvaluator.initialize(function, argumentInspectors)
checkCompatibleHiveReturnType(inspector, catalystReturnType)
Comment thread
srielau marked this conversation as resolved.
inspector
}

@transient
private lazy val deferredObjects: Array[DeferredObject] = argumentInspectors.zip(children).map {
case (inspect, child) => new DeferredObjectAdapter(inspect, child.dataType)
}

@transient
private lazy val unwrapper: Any => Any = unwrapperFor(returnInspector)
private lazy val unwrapper: Any => Any = unwrapperFor(returnInspector, catalystReturnType)

override def returnType: DataType = inspectorToDataType(returnInspector)
override def returnType: DataType = catalystReturnType

def setArg(index: Int, arg: Any): Unit =
deferredObjects(index).asInstanceOf[DeferredObjectAdapter].set(() => arg)
Expand Down
Loading