From e145b304ee7779ef1e66032429539444e0a4f42b Mon Sep 17 00:00:00 2001 From: srielau Date: Mon, 21 Sep 2026 05:15:10 +0000 Subject: [PATCH 1/3] feat: [SPARK-59683][SQL][HIVE] Support first-class CHAR/VARCHAR in script TRANSFORM output Support CHAR/VARCHAR output from script TRANSFORM with and without Hive SerDe when first-class semantics are enabled. Separated from SPARK-59277 (Hive inspector conversion and UDF/UDAF/UDTF support) per review feedback. JIRA: https://issues.apache.org/jira/browse/SPARK-59683 - No-SerDe: CHAR/VARCHAR scalar, nested, and complex output with padding, overflow, null-token, and malformed-field handling. - No-SerDe nested maps: restore JSON string keys to the declared physical key type through a per-row ArrayBasedMapBuilder that validates null and duplicate converted keys. - SerDe: rewrite LazySimpleSerDe CHAR/VARCHAR output types to STRING so Spark applies first-class length checks; non-LazySimpleSerDe SerDes keep the declared schema. --- .../BaseScriptTransformationExec.scala | 148 +++++++++++- .../BaseScriptTransformationSuite.scala | 211 ++++++++++++++++++ .../HiveScriptTransformationExec.scala | 26 ++- .../HiveScriptTransformationSuite.scala | 167 ++++++++++++++ 4 files changed, 545 insertions(+), 7 deletions(-) diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/BaseScriptTransformationExec.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/BaseScriptTransformationExec.scala index a313e2c671bec..fc47a2e21b944 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/BaseScriptTransformationExec.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/BaseScriptTransformationExec.scala @@ -31,10 +31,27 @@ import org.apache.spark.internal.Logging import org.apache.spark.internal.LogKeys._ import org.apache.spark.rdd.RDD import org.apache.spark.sql.catalyst.{CatalystTypeConverters, InternalRow} -import org.apache.spark.sql.catalyst.expressions.{Attribute, AttributeSet, Cast, Expression, GenericInternalRow, JsonToStructs, Literal, StructsToJson, UnsafeProjection} +import org.apache.spark.sql.catalyst.expressions.{ + Attribute, + AttributeSet, + BoundReference, + Cast, + Expression, + GenericInternalRow, + JsonToStructs, + Literal, + StructsToJson, + UnsafeProjection} import org.apache.spark.sql.catalyst.plans.logical.ScriptInputOutputSchema import org.apache.spark.sql.catalyst.plans.physical.Partitioning -import org.apache.spark.sql.catalyst.util.{DateTimeUtils, IntervalUtils} +import org.apache.spark.sql.catalyst.util.{ + ArrayBasedMapBuilder, + ArrayData, + CharVarcharUtils, + DateTimeUtils, + GenericArrayData, + IntervalUtils, + MapData} import org.apache.spark.sql.errors.QueryExecutionErrors import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.types._ @@ -201,7 +218,15 @@ trait BaseScriptTransformationExec extends UnaryExecNode { private lazy val outputFieldWriters: Seq[String => Any] = output.map { attr => val converter = CatalystTypeConverters.createToCatalystConverter(attr.dataType) attr.dataType match { - case StringType => wrapperConvertException(data => data, converter) + case _: CharType | _: VarcharType => + // First-class CHAR/VARCHAR must not use Hive LazySimpleSerDe's null-on-error path. + (data: String) => + if (data == ioschema.outputRowFormatMap("TOK_TABLEROWFORMATNULL")) { + null + } else { + converter(data) + } + case _: StringType => wrapperConvertException(data => data, converter) case BooleanType => wrapperConvertException(data => data.toBoolean, converter) case ByteType => wrapperConvertException(data => data.toByte, converter) case BinaryType => @@ -241,6 +266,29 @@ trait BaseScriptTransformationExec extends UnaryExecNode { data => IntervalUtils.microsToDuration( IntervalUtils.castStringToDTInterval(UTF8String.fromString(data), start, end)), converter) + case dt @ (_: ArrayType | _: MapType | _: StructType) + if CharVarcharUtils.hasCharVarchar(dt) => + val physicalType = ScriptTransformationIOSchema.toUnboundedStringType(dt) + // JSON object keys are strings. Cast them to the declared map key type after parsing. + val jsonType = ScriptTransformationIOSchema.toJsonMapKeyType(physicalType) + val complexTypeFactory = JsonToStructs( + jsonType, + ioschema.outputSerdeProps.toMap, + Literal(null), + Some(conf.sessionLocalTimeZone)) + val parsedToPhysical = if (jsonType.sameType(physicalType)) { + identity[Any] _ + } else { + val restoreMapKeys = ScriptTransformationIOSchema.makeJsonMapKeyRestorer( + physicalType, Some(conf.sessionLocalTimeZone)) + value: Any => restoreMapKeys(value) + } + val toScala = CatalystTypeConverters.createToScalaConverter(physicalType) + val parser = wrapperConvertException( + data => parsedToPhysical( + complexTypeFactory.nullSafeEval(UTF8String.fromString(data))), + identity) + data => converter(toScala(parser(data))) case _: ArrayType | _: MapType | _: StructType => val complexTypeFactory = JsonToStructs(attr.dataType, ioschema.outputSerdeProps.toMap, Literal(null), Some(conf.sessionLocalTimeZone)) @@ -253,7 +301,7 @@ trait BaseScriptTransformationExec extends UnaryExecNode { } } - // Keep consistent with Hive `LazySimpleSerde`, when there is a type case error, return null + // Match Hive `LazySimpleSerDe`: return null when a type cast fails. private val wrapperConvertException: (String => Any, Any => Any) => String => Any = (f: String => Any, converter: Any => Any) => (data: String) => converter { @@ -376,6 +424,98 @@ case class ScriptTransformationIOSchema( } object ScriptTransformationIOSchema { + private[sql] def toUnboundedStringType(dataType: DataType): DataType = { + dataType.transformRecursively { + case c: CharType => c.toStringType + case v: VarcharType => v.toStringType + } + } + + // JSON object keys are always strings. Rewrite every map key, including nested maps. + // `transformRecursively` would stop at the first matching MapType and skip children. + private[sql] def toJsonMapKeyType(dataType: DataType): DataType = dataType match { + case ArrayType(et, n) => ArrayType(toJsonMapKeyType(et), n) + case MapType(kt, vt, n) => + val jsonKey = if (kt.isInstanceOf[StringType]) kt else StringType + MapType(jsonKey, toJsonMapKeyType(vt), n) + case StructType(fields) => + StructType(fields.map(f => f.copy(dataType = toJsonMapKeyType(f.dataType)))) + case other => other + } + + /** + * Build a per-call map-key restorer that converts parsed JSON string keys + * back to the declared physical key type and validates the result through a + * fresh [[ArrayBasedMapBuilder]] on every invocation, so a failed key + * conversion or duplicate key cannot leave shared state dirty for the next row. + */ + private[sql] def makeJsonMapKeyRestorer( + targetType: DataType, + timeZoneId: Option[String]): Any => Any = { + val jsonType = toJsonMapKeyType(targetType) + + def make(jt: DataType, tt: DataType): Any => Any = (jt, tt) match { + case (ArrayType(jet, _), ArrayType(tet, _)) => + val elem = make(jet, tet) + (input: Any) => { + val arr = input.asInstanceOf[ArrayData] + val n = arr.numElements() + val out = new Array[Any](n) + var i = 0 + while (i < n) { + out(i) = if (arr.isNullAt(i)) null + else elem(arr.get(i, jet)) + i += 1 + } + new GenericArrayData(out) + } + + case (MapType(jkt, jvt, _), MapType(tkt, tvt, _)) => + val keyCast: Any => Any = if (jkt.sameType(tkt)) identity + else { + val c = Cast(BoundReference(0, jkt, nullable = false), + tkt, timeZoneId) + (k: Any) => c.eval(InternalRow(k)) + } + val valRestore = make(jvt, tvt) + (input: Any) => { + val map = input.asInstanceOf[MapData] + val n = map.numElements() + val builder = new ArrayBasedMapBuilder(tkt, tvt) + var i = 0 + while (i < n) { + val k = keyCast(map.keyArray().get(i, jkt)) + val v = if (map.valueArray().isNullAt(i)) null + else valRestore(map.valueArray().get(i, jvt)) + builder.put(k, v) + i += 1 + } + builder.build() + } + + case (js: StructType, ts: StructType) => + val restorers = js.fields.zip(ts.fields).map { + case (jf, tf) => make(jf.dataType, tf.dataType) + } + (input: Any) => { + val row = input.asInstanceOf[InternalRow] + val out = new GenericInternalRow(ts.length) + var i = 0 + while (i < ts.length) { + if (row.isNullAt(i)) out.setNullAt(i) + else out.update(i, + restorers(i)(row.get(i, js(i).dataType))) + i += 1 + } + out + } + + case _ => identity + } + + make(jsonType, targetType) + } + val defaultFormat = Map( ("TOK_TABLEROWFORMATFIELD", "\u0001"), ("TOK_TABLEROWFORMATLINES", "\n"), diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/BaseScriptTransformationSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/BaseScriptTransformationSuite.scala index 62c5f5631776b..d9f2902472087 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/BaseScriptTransformationSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/BaseScriptTransformationSuite.scala @@ -86,6 +86,217 @@ abstract class BaseScriptTransformationSuite extends QueryTest { assert(uncaughtExceptionHandler.exception.isEmpty) } + test("SPARK-59277: TRANSFORM output supports first-class CHAR/VARCHAR without SerDe") { + assume(TestUtils.testCommandAvailable("/bin/bash")) + withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true") { + val input = Seq(("ab", "xyz")).toDF("c", "v") + checkAnswer( + input, + (child: SparkPlan) => createScriptTransformationExec( + script = "cat", + output = Seq( + AttributeReference("c", CharType(4, "UTF8_LCASE"))(), + AttributeReference("v", VarcharType(5, "UNICODE_CI"))()), + child = child, + ioschema = defaultIOSchema), + Seq(Row("ab ", "xyz"))) + } + assert(uncaughtExceptionHandler.exception.isEmpty) + } + + test("SPARK-59277: TRANSFORM CHAR overflow without SerDe raises EXCEED_LIMIT_LENGTH") { + assume(TestUtils.testCommandAvailable("/bin/bash")) + withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true") { + val input = Seq("abcdef").toDF("c") + val exception = intercept[Exception] { + QueryTest.executePlan( + createScriptTransformationExec( + script = "cat", + output = Seq(AttributeReference("c", CharType(4))()), + child = input.queryExecution.sparkPlan, + ioschema = defaultIOSchema), + spark.sqlContext) + } + val runtimeException = exception match { + case s: org.apache.spark.SparkRuntimeException => s + case other => + other.getCause.asInstanceOf[org.apache.spark.SparkRuntimeException] + } + checkError( + exception = runtimeException, + condition = "EXCEED_LIMIT_LENGTH", + parameters = Map("limit" -> "4")) + } + } + + test("SPARK-59277: TRANSFORM VARCHAR overflow without SerDe raises EXCEED_LIMIT_LENGTH") { + assume(TestUtils.testCommandAvailable("/bin/bash")) + withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true") { + val input = Seq("abcdefgh").toDF("v") + val exception = intercept[Exception] { + QueryTest.executePlan( + createScriptTransformationExec( + script = "cat", + output = Seq(AttributeReference("v", VarcharType(5))()), + child = input.queryExecution.sparkPlan, + ioschema = defaultIOSchema), + spark.sqlContext) + } + val runtimeException = exception match { + case s: org.apache.spark.SparkRuntimeException => s + case other => + other.getCause.asInstanceOf[org.apache.spark.SparkRuntimeException] + } + checkError( + exception = runtimeException, + condition = "EXCEED_LIMIT_LENGTH", + parameters = Map("limit" -> "5")) + } + } + + test("SPARK-59277: TRANSFORM converts nested CHAR/VARCHAR without SerDe") { + assume(TestUtils.testCommandAvailable("/bin/bash")) + withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true") { + Seq( + ("""["ab"]""", ArrayType(CharType(4)), Row(Seq("ab "))), + ("""["xy"]""", ArrayType(VarcharType(4)), Row(Seq("xy"))), + ( + """{"1":"ab"}""", + MapType(IntegerType, CharType(4)), + Row(Map(1 -> "ab "))), + ( + """{"1":{"2":"ab"}}""", + MapType(IntegerType, MapType(IntegerType, CharType(4))), + Row(Map(1 -> Map(2 -> "ab ")))), + ( + """[{"1":"ab"}]""", + ArrayType(MapType(IntegerType, CharType(4))), + Row(Seq(Map(1 -> "ab ")))), + ( + """{"m":{"1":"ab"}}""", + StructType(Seq(StructField("m", MapType(IntegerType, CharType(4))))), + Row(Row(Map(1 -> "ab ")))), + ( + """{"value":"xy"}""", + StructType(Seq(StructField("value", CharType(5)))), + Row(Row("xy ")))).foreach { case (json, dataType, expected) => + val input = Seq(json).toDF("value") + checkAnswer( + input, + (child: SparkPlan) => createScriptTransformationExec( + script = "cat", + output = Seq(AttributeReference("value", dataType)()), + child = child, + ioschema = defaultIOSchema), + Seq(expected)) + } + } + assert(uncaughtExceptionHandler.exception.isEmpty) + } + + test("SPARK-59277: TRANSFORM nested CHAR/VARCHAR overflow without SerDe") { + assume(TestUtils.testCommandAvailable("/bin/bash")) + withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true") { + Seq( + (ArrayType(CharType(4)), """["abcdef"]"""), + (MapType(IntegerType, CharType(4)), """{"1":"abcdef"}"""), + ( + MapType(IntegerType, MapType(IntegerType, CharType(4))), + """{"1":{"2":"abcdef"}}"""), + ( + StructType(Seq(StructField("value", VarcharType(4)))), + """{"value":"abcdef"}""")).foreach { case (dataType, json) => + val input = Seq(json).toDF("value") + val exception = intercept[Exception] { + QueryTest.executePlan( + createScriptTransformationExec( + script = "cat", + output = Seq(AttributeReference("value", dataType)()), + child = input.queryExecution.sparkPlan, + ioschema = defaultIOSchema), + spark.sqlContext) + } + val runtimeException = exception match { + case s: org.apache.spark.SparkRuntimeException => s + case other => + other.getCause.asInstanceOf[org.apache.spark.SparkRuntimeException] + } + checkError( + exception = runtimeException, + condition = "EXCEED_LIMIT_LENGTH", + parameters = Map("limit" -> "4")) + } + } + } + + test("SPARK-59277: malformed nested CHAR JSON without SerDe returns null") { + assume(TestUtils.testCommandAvailable("/bin/bash")) + withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true") { + val input = Seq("""{"1":""").toDF("value") + checkAnswer( + input, + (child: SparkPlan) => createScriptTransformationExec( + script = "cat", + output = Seq( + AttributeReference("value", MapType(IntegerType, CharType(4)))()), + child = child, + ioschema = defaultIOSchema), + Seq(Row(null))) + } + assert(uncaughtExceptionHandler.exception.isEmpty) + } + + test("SPARK-59277: TRANSFORM validates restored JSON map keys without SerDe") { + assume(TestUtils.testCommandAvailable("/bin/bash")) + withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true") { + val mapType = MapType(IntegerType, CharType(4)) + Seq( + ("""{"1":"ab"}""", mapType, Row(Map(1 -> "ab "))), + ("""{"not-an-int":"ab"}""", mapType, Row(null)), + ("""{"1":"a","01":"b"}""", mapType, Row(null)), + ( + """[{"not-an-int":"ab"}]""", + ArrayType(mapType), + Row(null)), + ( + """{"m":{"1":"a","01":"b"}}""", + StructType(Seq(StructField("m", mapType))), + Row(null))).foreach { case (json, dataType, expected) => + val input = Seq(json).toDF("value") + checkAnswer( + input, + (child: SparkPlan) => createScriptTransformationExec( + script = "cat", + output = Seq(AttributeReference("value", dataType)()), + child = child, + ioschema = defaultIOSchema), + Seq(expected)) + } + } + assert(uncaughtExceptionHandler.exception.isEmpty) + } + + test("SPARK-59277: colliding map key followed by valid row without SerDe") { + assume(TestUtils.testCommandAvailable("/bin/bash")) + withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true") { + val mapType = MapType(IntegerType, CharType(4)) + // Row 1 has duplicate converted keys (1 and 01 both cast to 1). + // Row 2 is valid. Both rows are in the same partition. + val input = Seq( + """{"1":"a","01":"b"}""", + """{"2":"cd"}""").toDF("value").coalesce(1) + checkAnswer( + input, + (child: SparkPlan) => createScriptTransformationExec( + script = "cat", + output = Seq(AttributeReference("value", mapType)()), + child = child, + ioschema = defaultIOSchema), + Seq(Row(null), Row(Map(2 -> "cd ")))) + } + assert(uncaughtExceptionHandler.exception.isEmpty) + } + test("script transformation should not swallow errors from upstream operators (no serde)") { assume(TestUtils.testCommandAvailable("/bin/bash")) diff --git a/sql/hive/src/main/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationExec.scala b/sql/hive/src/main/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationExec.scala index de2d15415837a..107530f8ffbce 100644 --- a/sql/hive/src/main/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationExec.scala +++ b/sql/hive/src/main/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationExec.scala @@ -27,6 +27,7 @@ import org.apache.hadoop.conf.Configuration import org.apache.hadoop.hive.ql.exec.{RecordReader, RecordWriter} import org.apache.hadoop.hive.serde.serdeConstants import org.apache.hadoop.hive.serde2.AbstractSerDe +import org.apache.hadoop.hive.serde2.`lazy`.LazySimpleSerDe import org.apache.hadoop.hive.serde2.objectinspector._ import org.apache.hadoop.io.Writable @@ -36,7 +37,7 @@ import org.apache.spark.sql.catalyst.expressions._ import org.apache.spark.sql.execution._ import org.apache.spark.sql.hive.HiveInspectors import org.apache.spark.sql.hive.HiveShim._ -import org.apache.spark.sql.types.DataType +import org.apache.spark.sql.types._ import org.apache.spark.util.{CircularBuffer, Utils} /** @@ -75,7 +76,9 @@ private[hive] case class HiveScriptTransformationExec( val mutableRow = new SpecificInternalRow(output.map(_.dataType)) @transient - lazy val unwrappers = outputSoi.getAllStructFieldRefs.asScala.map(unwrapperFor) + lazy val unwrappers = outputSoi.getAllStructFieldRefs.asScala.zip(output).map { + case (field, attr) => unwrapperFor(field, attr.dataType) + } override def hasNext: Boolean = { if (completed) { @@ -256,7 +259,8 @@ object HiveScriptIOSchema extends HiveInspectors { output: Seq[Attribute]): Option[(AbstractSerDe, StructObjectInspector)] = { ioschema.outputSerdeClass.map { serdeClass => val (columns, columnTypes) = parseAttrs(output) - val serde = initSerDe(serdeClass, columns, columnTypes, ioschema.outputSerdeProps) + val serdeTypes = outputTypesForSerDe(serdeClass, columnTypes) + val serde = initSerDe(serdeClass, columns, serdeTypes, ioschema.outputSerdeProps) val structObjectInspector = serde.getObjectInspector().asInstanceOf[StructObjectInspector] (serde, structObjectInspector) } @@ -268,6 +272,22 @@ object HiveScriptIOSchema extends HiveInspectors { (columns, columnTypes) } + /** + * Hive LazySimpleSerDe CHAR/VARCHAR types truncate on deserialize. Map them to STRING so + * Spark applies first-class length checks. Only LazySimpleSerDe and subclasses get this + * rewrite; other SerDes keep the declared CHAR/VARCHAR schema. + */ + private def outputTypesForSerDe( + serdeClassName: String, + columnTypes: Seq[DataType]): Seq[DataType] = { + val serdeClass = Utils.classForName[AbstractSerDe](serdeClassName) + if (classOf[LazySimpleSerDe].isAssignableFrom(serdeClass)) { + columnTypes.map(ScriptTransformationIOSchema.toUnboundedStringType) + } else { + columnTypes + } + } + def initSerDe( serdeClassName: String, columns: Seq[String], diff --git a/sql/hive/src/test/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationSuite.scala b/sql/hive/src/test/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationSuite.scala index b1ff05b8c1b06..5772784795fd1 100644 --- a/sql/hive/src/test/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationSuite.scala +++ b/sql/hive/src/test/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationSuite.scala @@ -20,8 +20,15 @@ package org.apache.spark.sql.hive.execution import java.sql.Timestamp import java.time.{Duration, Period} import java.time.temporal.ChronoUnit +import java.util.{Arrays, Properties} +import org.apache.hadoop.conf.Configuration +import org.apache.hadoop.hive.serde.serdeConstants +import org.apache.hadoop.hive.serde2.{AbstractSerDe, SerDeStats} import org.apache.hadoop.hive.serde2.`lazy`.LazySimpleSerDe +import org.apache.hadoop.hive.serde2.objectinspector.{ObjectInspector, ObjectInspectorFactory} +import org.apache.hadoop.hive.serde2.objectinspector.primitive.PrimitiveObjectInspectorFactory +import org.apache.hadoop.io.{Text, Writable} import org.scalatest.exceptions.TestFailedException import org.apache.spark.{SparkException, TestUtils} @@ -31,6 +38,7 @@ import org.apache.spark.sql.catalyst.util.DateTimeConstants import org.apache.spark.sql.execution._ import org.apache.spark.sql.functions._ import org.apache.spark.sql.hive.test.TestHiveSingleton +import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.types._ import org.apache.spark.sql.types.DayTimeIntervalType._ import org.apache.spark.sql.types.YearMonthIntervalType._ @@ -372,6 +380,139 @@ class HiveScriptTransformationSuite extends BaseScriptTransformationSuite with T } } + test("SPARK-59277: TRANSFORM supports nested collated CHAR/VARCHAR with Hive SerDe") { + assume(TestUtils.testCommandAvailable("/bin/bash")) + withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true") { + val query = sql( + """ + |SELECT TRANSFORM( + | array(CAST('ab' AS CHAR(4) COLLATE UTF8_LCASE)), + | named_struct('value', CAST('xyz' AS VARCHAR(6) COLLATE UNICODE_CI))) + |USING 'cat' + |AS ( + | chars ARRAY, + | nested STRUCT) + |FROM VALUES (1) input(dummy) + |""".stripMargin) + assert(query.schema.map(_.dataType) === Seq( + ArrayType(CharType(4, "UTF8_LCASE")), + StructType(Seq(StructField("value", VarcharType(6, "UNICODE_CI")))))) + checkAnswer(query, Row(Seq("ab "), Row("xyz"))) + } + } + + test("SPARK-59277: TRANSFORM CHAR overflow with Hive SerDe raises EXCEED_LIMIT_LENGTH") { + assume(TestUtils.testCommandAvailable("/bin/bash")) + withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true") { + val exception = intercept[Exception] { + sql( + """ + |SELECT TRANSFORM('abcdef') + |USING 'cat' + |AS (c CHAR(4)) + |FROM VALUES (1) input(dummy) + |""".stripMargin).collect() + } + val runtimeException = exception match { + case s: org.apache.spark.SparkRuntimeException => s + case other => + other.getCause.asInstanceOf[org.apache.spark.SparkRuntimeException] + } + checkError( + exception = runtimeException, + condition = "EXCEED_LIMIT_LENGTH", + parameters = Map("limit" -> "4")) + } + } + + test("SPARK-59277: TRANSFORM VARCHAR overflow with Hive SerDe raises EXCEED_LIMIT_LENGTH") { + assume(TestUtils.testCommandAvailable("/bin/bash")) + withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true") { + val exception = intercept[Exception] { + sql( + """ + |SELECT TRANSFORM('abcdefgh') + |USING 'cat' + |AS (v VARCHAR(5)) + |FROM VALUES (1) input(dummy) + |""".stripMargin).collect() + } + val runtimeException = exception match { + case s: org.apache.spark.SparkRuntimeException => s + case other => + other.getCause.asInstanceOf[org.apache.spark.SparkRuntimeException] + } + checkError( + exception = runtimeException, + condition = "EXCEED_LIMIT_LENGTH", + parameters = Map("limit" -> "5")) + } + } + + test("SPARK-59277: nested CHAR/VARCHAR overflow with Hive SerDe") { + assume(TestUtils.testCommandAvailable("/bin/bash")) + withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true") { + Seq( + """ + |SELECT TRANSFORM(array('abcdef')) + |USING 'cat' + |AS (value ARRAY) + |FROM VALUES (1) input(dummy) + |""".stripMargin, + """ + |SELECT TRANSFORM(named_struct('value', 'abcdef')) + |USING 'cat' + |AS (value STRUCT) + |FROM VALUES (1) input(dummy) + |""".stripMargin).foreach { query => + val exception = intercept[Exception] { + sql(query).collect() + } + val runtimeException = exception match { + case s: org.apache.spark.SparkRuntimeException => s + case other => + other.getCause.asInstanceOf[org.apache.spark.SparkRuntimeException] + } + checkError( + exception = runtimeException, + condition = "EXCEED_LIMIT_LENGTH", + parameters = Map("limit" -> "4")) + } + } + } + + test("SPARK-59277: output SerDe CHAR/VARCHAR rewrite is LazySimpleSerDe-only") { + withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true") { + val output = Seq( + AttributeReference("c", CharType(4))(), + AttributeReference("v", VarcharType(5))(), + AttributeReference("nested", ArrayType(CharType(4)))()) + + val (_, lazySoi) = HiveScriptIOSchema.initOutputSerDe(hiveIOSchema, output).get + assert(lazySoi.getAllStructFieldRefs.get(0).getFieldObjectInspector.getTypeName === + "string") + assert(lazySoi.getAllStructFieldRefs.get(1).getFieldObjectInspector.getTypeName === + "string") + assert(lazySoi.getAllStructFieldRefs.get(2).getFieldObjectInspector.getTypeName === + "array") + + val subclassSchema = hiveIOSchema.copy( + outputSerdeClass = Some(classOf[TestLazySimpleSerDe].getCanonicalName)) + val (_, subclassSoi) = HiveScriptIOSchema.initOutputSerDe(subclassSchema, output).get + assert(subclassSoi.getAllStructFieldRefs.get(0).getFieldObjectInspector.getTypeName === + "string") + + SchemaCapturingSerDe.lastColumnTypes = null + val customSchema = defaultIOSchema.copy( + outputSerdeClass = Some(classOf[SchemaCapturingSerDe].getCanonicalName)) + HiveScriptIOSchema.initOutputSerDe(customSchema, output) + val captured = SchemaCapturingSerDe.lastColumnTypes + assert(captured.contains("char(4)")) + assert(captured.contains("varchar(5)")) + assert(captured.contains("array")) + } + } + test("SPARK-32400: TRANSFORM doesn't support CalendarIntervalType/UserDefinedType (hive serde)") { assume(TestUtils.testCommandAvailable("/bin/bash")) withTempView("v") { @@ -668,3 +809,29 @@ class HiveScriptTransformationSuite extends BaseScriptTransformationSuite with T } } } + +class TestLazySimpleSerDe extends LazySimpleSerDe + +class SchemaCapturingSerDe extends AbstractSerDe { + override def initialize(conf: Configuration, tbl: Properties): Unit = { + SchemaCapturingSerDe.lastColumnTypes = + tbl.getProperty(serdeConstants.LIST_COLUMN_TYPES) + } + + override def getObjectInspector: ObjectInspector = + ObjectInspectorFactory.getStandardStructObjectInspector( + Arrays.asList("col"), + Arrays.asList(PrimitiveObjectInspectorFactory.javaStringObjectInspector)) + + override def getSerializedClass: Class[_ <: Writable] = classOf[Text] + + override def getSerDeStats: SerDeStats = null + + override def serialize(obj: Any, inspector: ObjectInspector): Writable = null + + override def deserialize(blob: Writable): AnyRef = null +} + +object SchemaCapturingSerDe { + @volatile var lastColumnTypes: String = _ +} From 051fa0c17e42110be6f311ea3d3976595bd1e670 Mon Sep 17 00:00:00 2001 From: srielau Date: Wed, 30 Sep 2026 17:22:17 +0000 Subject: [PATCH 2/3] [SPARK-59683][SQL][HIVE] Gate TRANSFORM CHAR/VARCHAR on standard semantics preserveCharVarcharTypeInfo predates this work and must keep the old unsupported/no-rewrite TRANSFORM path when standard semantics is off. --- .../BaseScriptTransformationExec.scala | 16 ++++++++-- .../BaseScriptTransformationSuite.scala | 30 +++++++++++++++++++ .../HiveScriptTransformationExec.scala | 11 ++++--- .../HiveScriptTransformationSuite.scala | 24 +++++++++++++++ 4 files changed, 74 insertions(+), 7 deletions(-) diff --git a/sql/core/src/main/scala/org/apache/spark/sql/execution/BaseScriptTransformationExec.scala b/sql/core/src/main/scala/org/apache/spark/sql/execution/BaseScriptTransformationExec.scala index fc47a2e21b944..2d24349568a32 100644 --- a/sql/core/src/main/scala/org/apache/spark/sql/execution/BaseScriptTransformationExec.scala +++ b/sql/core/src/main/scala/org/apache/spark/sql/execution/BaseScriptTransformationExec.scala @@ -79,6 +79,10 @@ trait BaseScriptTransformationExec extends UnaryExecNode { override def outputPartitioning: Partitioning = child.outputPartitioning + // Snapshot on the driver. SparkPlan.session is @transient, so executor `conf` is the + // default SQLConf and would ignore withSQLConf / session standard-semantics. + private val standardCharVarcharSemantics: Boolean = conf.charVarcharStandardSemantics + override def doExecute(): RDD[InternalRow] = { val broadcastedHadoopConf = new SerializableConfiguration(session.sessionState.newHadoopConf()) @@ -218,8 +222,14 @@ trait BaseScriptTransformationExec extends UnaryExecNode { private lazy val outputFieldWriters: Seq[String => Any] = output.map { attr => val converter = CatalystTypeConverters.createToCatalystConverter(attr.dataType) attr.dataType match { - case _: CharType | _: VarcharType => - // First-class CHAR/VARCHAR must not use Hive LazySimpleSerDe's null-on-error path. + case dt @ (_: CharType | _: VarcharType) => + // Preserve-only keeps CHAR/VARCHAR in the schema but does not apply SQL pad/overflow. + // Match CHAR/VARCHAR before `_: StringType` (they extend StringType). + if (!standardCharVarcharSemantics) { + throw QueryExecutionErrors.outputDataTypeUnsupportedByNodeWithoutSerdeError( + nodeName, dt) + } + // Do not use the SerDe null-on-error wrapper; that would hide EXCEED_LIMIT_LENGTH. (data: String) => if (data == ioschema.outputRowFormatMap("TOK_TABLEROWFORMATNULL")) { null @@ -267,7 +277,7 @@ trait BaseScriptTransformationExec extends UnaryExecNode { IntervalUtils.castStringToDTInterval(UTF8String.fromString(data), start, end)), converter) case dt @ (_: ArrayType | _: MapType | _: StructType) - if CharVarcharUtils.hasCharVarchar(dt) => + if standardCharVarcharSemantics && CharVarcharUtils.hasCharVarchar(dt) => val physicalType = ScriptTransformationIOSchema.toUnboundedStringType(dt) // JSON object keys are strings. Cast them to the declared map key type after parsing. val jsonType = ScriptTransformationIOSchema.toJsonMapKeyType(physicalType) diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/BaseScriptTransformationSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/BaseScriptTransformationSuite.scala index d9f2902472087..1d882a21d42d9 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/BaseScriptTransformationSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/BaseScriptTransformationSuite.scala @@ -297,6 +297,36 @@ abstract class BaseScriptTransformationSuite extends QueryTest { assert(uncaughtExceptionHandler.exception.isEmpty) } + test("SPARK-59683: preserve-only CHAR/VARCHAR without SerDe stays unsupported") { + assume(TestUtils.testCommandAvailable("/bin/bash")) + withSQLConf( + SQLConf.PRESERVE_CHAR_VARCHAR_TYPE_INFO.key -> "true", + SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "false") { + val input = Seq("ab").toDF("c") + val exception = intercept[Exception] { + QueryTest.executePlan( + createScriptTransformationExec( + script = "cat", + output = Seq(AttributeReference("c", CharType(4))()), + child = input.queryExecution.sparkPlan, + ioschema = defaultIOSchema), + spark.sqlContext) + } + var cur: Throwable = exception + var sparkException: SparkException = null + while (cur != null && sparkException == null) { + cur match { + case s: SparkException => sparkException = s + case _ => + } + cur = cur.getCause + } + assert(sparkException != null, exception) + assert(sparkException.getCondition === "_LEGACY_ERROR_TEMP_2265") + assert(sparkException.getMessageParameters.get("dt") === "CharType") + } + } + test("script transformation should not swallow errors from upstream operators (no serde)") { assume(TestUtils.testCommandAvailable("/bin/bash")) diff --git a/sql/hive/src/main/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationExec.scala b/sql/hive/src/main/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationExec.scala index 107530f8ffbce..13ed2fb844af4 100644 --- a/sql/hive/src/main/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationExec.scala +++ b/sql/hive/src/main/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationExec.scala @@ -37,6 +37,7 @@ import org.apache.spark.sql.catalyst.expressions._ import org.apache.spark.sql.execution._ import org.apache.spark.sql.hive.HiveInspectors import org.apache.spark.sql.hive.HiveShim._ +import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.types._ import org.apache.spark.util.{CircularBuffer, Utils} @@ -273,15 +274,17 @@ object HiveScriptIOSchema extends HiveInspectors { } /** - * Hive LazySimpleSerDe CHAR/VARCHAR types truncate on deserialize. Map them to STRING so - * Spark applies first-class length checks. Only LazySimpleSerDe and subclasses get this - * rewrite; other SerDes keep the declared CHAR/VARCHAR schema. + * Hive LazySimpleSerDe CHAR/VARCHAR types truncate on deserialize. Under standard + * semantics, map them to STRING so Spark applies length checks. Only LazySimpleSerDe + * and subclasses get this rewrite; other SerDes keep the declared schema. + * preserveCharVarcharTypeInfo without standard semantics is unchanged. */ private def outputTypesForSerDe( serdeClassName: String, columnTypes: Seq[DataType]): Seq[DataType] = { val serdeClass = Utils.classForName[AbstractSerDe](serdeClassName) - if (classOf[LazySimpleSerDe].isAssignableFrom(serdeClass)) { + if (SQLConf.get.charVarcharStandardSemantics && + classOf[LazySimpleSerDe].isAssignableFrom(serdeClass)) { columnTypes.map(ScriptTransformationIOSchema.toUnboundedStringType) } else { columnTypes diff --git a/sql/hive/src/test/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationSuite.scala b/sql/hive/src/test/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationSuite.scala index 5772784795fd1..b969b37697320 100644 --- a/sql/hive/src/test/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationSuite.scala +++ b/sql/hive/src/test/scala/org/apache/spark/sql/hive/execution/HiveScriptTransformationSuite.scala @@ -513,6 +513,30 @@ class HiveScriptTransformationSuite extends BaseScriptTransformationSuite with T } } + test("SPARK-59683: preserve-only does not rewrite LazySimpleSerDe CHAR/VARCHAR") { + withSQLConf( + SQLConf.PRESERVE_CHAR_VARCHAR_TYPE_INFO.key -> "true", + SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "false") { + val output = Seq( + AttributeReference("c", CharType(4))(), + AttributeReference("v", VarcharType(5))(), + AttributeReference("nested", ArrayType(CharType(4)))()) + scala.util.Try(HiveScriptIOSchema.initOutputSerDe(hiveIOSchema, output).get) match { + case scala.util.Success((_, soi)) => + val typeName = soi.getAllStructFieldRefs.get(0).getFieldObjectInspector.getTypeName + assert(typeName === "char(4)", + s"preserve-only must not rewrite CHAR to STRING, found $typeName") + assert(soi.getAllStructFieldRefs.get(1).getFieldObjectInspector.getTypeName === + "varchar(5)") + assert(soi.getAllStructFieldRefs.get(2).getFieldObjectInspector.getTypeName === + "array") + case scala.util.Failure(_) => + // Hive TypeInfo may reject CHAR on this lineage. The standard-semantics + // rewrite would have succeeded with STRING inspectors. + } + } + } + test("SPARK-32400: TRANSFORM doesn't support CalendarIntervalType/UserDefinedType (hive serde)") { assume(TestUtils.testCommandAvailable("/bin/bash")) withTempView("v") { From bf2f8a0ba89284d968197d53da755ab0f19cfa86 Mon Sep 17 00:00:00 2001 From: srielau Date: Wed, 30 Sep 2026 17:56:58 +0000 Subject: [PATCH 3/3] [SPARK-59683][HIVE] Add DataType-aware in-place Hive field unwrapper TRANSFORM SerDe output needs unwrapperFor(field, dataType) so CHAR/VARCHAR length checks run after deserialize without dropping primitive setters. --- .../apache/spark/sql/hive/HiveInspectors.scala | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) diff --git a/sql/hive/src/main/scala/org/apache/spark/sql/hive/HiveInspectors.scala b/sql/hive/src/main/scala/org/apache/spark/sql/hive/HiveInspectors.scala index b8758407877f0..1eb4b48e8338e 100644 --- a/sql/hive/src/main/scala/org/apache/spark/sql/hive/HiveInspectors.scala +++ b/sql/hive/src/main/scala/org/apache/spark/sql/hive/HiveInspectors.scala @@ -942,6 +942,24 @@ private[hive] trait HiveInspectors { unwrapperFor(objectInspector) } + /** + * In-place unwrapper that uses DataType-aware conversion when the target + * contains CHAR/VARCHAR or nanosecond timestamps. Other fields keep the + * specialized primitive setters from the field-only overload. + */ + def unwrapperFor( + field: HiveStructField, + dataType: DataType): (Any, InternalRow, Int) => Unit = { + if (CharVarcharUtils.hasCharVarchar(dataType) || + dataType.existsRecursively(_.isInstanceOf[AnyTimestampNanoType])) { + val unwrapper = unwrapperFor(field.getFieldObjectInspector, dataType) + (value: Any, row: InternalRow, ordinal: Int) => + row.update(ordinal, unwrapper(value)) + } else { + unwrapperFor(field) + } + } + /** * Builds unwrappers ahead of time according to object inspector * types to avoid pattern matching and branching costs per row.