Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
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
7 changes: 7 additions & 0 deletions common/utils/src/main/resources/error/error-conditions.json
Original file line number Diff line number Diff line change
Expand Up @@ -9896,6 +9896,13 @@
],
"sqlState" : "0A000"
},
"UNSUPPORTED_XML_CHAR_VARCHAR_MAP_KEY" : {
"message" : [
"XML map keys of CHAR or VARCHAR cannot be padded or trimmed.",
"The key <key> is not valid for type <dataType>."
],
"sqlState" : "0A000"
},
"UNTYPED_SCALA_UDF" : {
"message" : [
"You're using untyped Scala UDF, which does not have the input type information. Spark may blindly pass null to the Scala closure with primitive-type argument, and the closure will see the default value of the Java type for the null argument, e.g. `udf((x: Int) => x, IntegerType)`, the result is 0 for null input. To get rid of this error, you could:",
Expand Down
1 change: 1 addition & 0 deletions docs/sql-migration-guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ license: |
## Upgrading from Spark SQL 4.3 to 4.4

- Since Spark 4.4, when `spark.sql.preserveCharVarcharTypeInfo` is true and `spark.sql.charVarchar.standardSemantics.enabled` is false, ORC reads that apply a CHAR/VARCHAR schema over STRING storage return the stored values without ORC truncation, matching Parquet. Previously the ORC reader requested `char(n)`/`varchar(n)` and truncated STRING-stored values to `n`. Read-side length checks (`EXCEED_LIMIT_LENGTH`) apply only when `spark.sql.charVarchar.standardSemantics.enabled` is true.
- Since Spark 4.4, when `spark.sql.charVarchar.standardSemantics.enabled` is true, XML names used as `MAP<CHAR(n), _>` or `MAP<VARCHAR(n), _>` keys in `from_xml` and the XML datasource are length-checked without padding or trimming. The checked names are element names, `attributePrefix` plus attribute names (default prefix `_`), and the `valueTag` (default `_VALUE`) when mixed text is present. A `CHAR(n)` key must already be exactly `n` characters, and a `VARCHAR(n)` key must already be at most `n` characters. Mismatched keys fail with `UNSUPPORTED_XML_CHAR_VARCHAR_MAP_KEY` (SQLSTATE `0A000`) and follow the XML parse mode (`PERMISSIVE` or `FAILFAST`). Exact repeated names last-win; `spark.sql.mapKeyDedupPolicy` is not applied. XML CHAR/VARCHAR values use the same pad and `EXCEED_LIMIT_LENGTH` checks as other parsed text. Empty, attribute-only, and text-only map elements stay SQL NULL or a malformed record; they do not become an empty map or a `valueTag` entry. In Spark 4.3, `from_xml` and XML reader schemas rejected CHAR/VARCHAR with `UNSUPPORTED_CHAR_OR_VARCHAR_AS_STRING`, or replaced them with STRING when `spark.sql.legacy.charVarcharAsString` was true, so CHAR/VARCHAR map keys did not reach the parser. When the standard-semantics flag is false and `spark.sql.preserveCharVarcharTypeInfo` is true, short CHAR keys are padded and over-length CHAR or VARCHAR keys raise `EXCEED_LIMIT_LENGTH`.
- Since Spark 4.4, the options maps passed to `from_csv`, `to_csv`, `schema_of_csv`, `from_json`, `to_json`, `schema_of_json`, `from_xml`, `to_xml`, and `schema_of_xml` must be foldable after replacing `RuntimeReplaceable` expressions. Previously, Spark evaluated non-foldable options during analysis, which allowed some constant expressions but could fail with an internal error or incorrectly evaluate row-dependent expressions. To allow deterministic and row-independent non-foldable options, set `spark.sql.legacy.allowNonFoldableOptions` to `true`. Row-dependent, unevaluable, and nondeterministic options are always rejected.
- Since Spark 4.4, when an already-analyzed Data Source V2 query is refreshed after a compatible schema change, connectors can return more data columns from the current table schema in `Scan.readSchema()` than requested by `SupportsPushDownRequiredColumns.pruneColumns`. Previously, this partial pruning could fail planning because the scan reported columns absent from the analyzed relation output.
- Since Spark 4.4, when an already-analyzed Data Source V2 query is refreshed, if rebinding an expanded struct used as a map key causes distinct current keys to collide under the analyzed schema, the query fails with `DUPLICATED_MAP_KEY` by default instead of returning duplicate keys. With `spark.sql.mapKeyDedupPolicy=LAST_WIN`, the last value is kept.
Expand Down

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,6 @@ import javax.xml.stream.events._
import javax.xml.transform.stream.StreamSource
import javax.xml.validation.Schema

import scala.collection.mutable
import scala.collection.mutable.ArrayBuffer
import scala.jdk.CollectionConverters._
import scala.util.Try
Expand All @@ -42,7 +41,7 @@ import org.apache.spark.{SparkIllegalArgumentException, SparkUpgradeException}
import org.apache.spark.internal.Logging
import org.apache.spark.sql.catalyst.InternalRow
import org.apache.spark.sql.catalyst.expressions.{ExprUtils, GenericInternalRow, ToStringBase}
import org.apache.spark.sql.catalyst.util.{ArrayBasedMapData, BadRecordException, CharVarcharUtils, DateFormatter, DropMalformedMode, DuplicateMapKeyUtils, FailureSafeParser, GenericArrayData, MapData, ParseMode, PartialResultArrayException, PartialResultException, PermissiveMode, TimeFormatter, TimestampFormatter}
import org.apache.spark.sql.catalyst.util.{ArrayBasedMapData, BadRecordException, CharVarcharUtils, DateFormatter, DropMalformedMode, FailureSafeParser, GenericArrayData, MapData, ParseMode, PartialResultArrayException, PartialResultException, PermissiveMode, TimeFormatter, TimestampFormatter}
import org.apache.spark.sql.catalyst.util.LegacyDateFormats.FAST_DATE_FORMAT
import org.apache.spark.sql.catalyst.xml.StaxXmlParser.convertStream
import org.apache.spark.sql.errors.QueryExecutionErrors
Expand Down Expand Up @@ -86,6 +85,9 @@ class StaxXmlParser(

private val caseSensitive = SQLConf.get.caseSensitiveAnalysis

// CHAR/VARCHAR XML map keys are length-checked without pad or trim under this flag.
private val charVarcharStandardSemantics = SQLConf.get.charVarcharStandardSemantics

/**
* Limits a view of an event stream to one element whose start event has already been consumed.
* Closing or draining this view consumes the matching end event without closing the underlying
Expand Down Expand Up @@ -233,7 +235,6 @@ class StaxXmlParser(
// ValidatorUtil.newValidator throws this when the JAXP implementation cannot
// disable external access; that is an environment error, not a bad record.
case e: UnsupportedOperationException => throw e
case DuplicateMapKeyUtils(e) => throw e
case e@(_: RuntimeException | _: XMLStreamException | _: MalformedInputException
| _: SAXException) =>
// XML parser currently doesn't support partial results for corrupted records.
Expand Down Expand Up @@ -330,16 +331,11 @@ class StaxXmlParser(
throw BadRecordException(xmlLiteral, () => Array.empty,
wrappedCharException)
case PartialResultException(row, cause) =>
DuplicateMapKeyUtils.cause(cause) match {
case Some(e) => throw e
case None =>
throw BadRecordException(record = xmlLiteral, partialResults = () => Array(row), cause)
}
throw BadRecordException(record = xmlLiteral, partialResults = () => Array(row), cause)
case PartialResultArrayException(rows, cause) =>
throw BadRecordException(record = xmlLiteral, partialResults = () => rows, cause)
case e: Throwable =>
SparkErrorUtils.getRootCause(e) match {
case DuplicateMapKeyUtils(duplicate) => throw duplicate
case _: FileNotFoundException if options.ignoreMissingFiles =>
logWarning("Skipped missing file", e)
parser.close()
Expand Down Expand Up @@ -383,9 +379,10 @@ class StaxXmlParser(
startElementName: String,
attributes: Array[Attribute]): Any = dt match {
case st: StructType => convertObject(parser, st)
case MapType(StringType, vt, _) => convertMap(parser, vt, attributes)
case MapType(kt @ (_: CharType | _: VarcharType), vt, _) =>
convertConstrainedMap(parser, kt, vt, attributes)
// CHAR/VARCHAR extend StringType. Non-string keys (for example MAP<INT, _> through
// format("xml").load()) must not reach convertMap: applyTextParseSemantics would keep
// UTF8String keys and fail later with ClassCastException.
case MapType(kt: StringType, vt, _) => convertMap(parser, vt, attributes, kt)
case ArrayType(st, _) => convertField(parser, st, startElementName)
case VariantType =>
StaxXmlParser.convertVariant(parser, attributes, options)
Expand Down Expand Up @@ -445,109 +442,80 @@ class StaxXmlParser(
}

/**
* Parse an object as map.
* Parse an object as a Map.
*
* When `spark.sql.charVarchar.standardSemantics.enabled` is true, XML names used as
* CHAR/VARCHAR keys are length-checked without rewriting: CHAR keys must already be
* exactly n characters, and VARCHAR keys must already be at most n characters.
* Padding, trimming, and mapKeyDedupPolicy are not applied. With the flag off,
* short CHAR keys are padded and over-length CHAR or VARCHAR keys raise
* EXCEED_LIMIT_LENGTH when first-class CHAR/VARCHAR types reach the parser.
* Repeated names last-win on binary equality (collation is not consulted), matching
* ordinary MAP<STRING, ...> XML maps.
*
* Example (flag on): `from_xml('<ROW><m><ab>1</ab></m></ROW>',
* 'm MAP<CHAR(2), INT>')` keeps key `ab`; `MAP<CHAR(4), INT>` raises
* `UNSUPPORTED_XML_CHAR_VARCHAR_MAP_KEY`.
*
* This method owns element bounding so a key-check failure still consumes the current
* map element (ARRAY<MAP<...>> and nested maps included).
*/
private def convertMap(
parser: XMLEventReader,
valueType: DataType,
attributes: Array[Attribute]): MapData = {
val kvPairs = ArrayBuffer.empty[(UTF8String, Any)]
attributes.foreach { attr =>
kvPairs += (UTF8String.fromString(options.attributePrefix + attr.getName.getLocalPart)
-> convertTo(attr.getValue, valueType))
}
var shouldStop = false
while (!shouldStop) {
parser.nextEvent match {
case e: StartElement =>
val key = StaxXmlParserUtils.getName(e.asStartElement.getName, options)
kvPairs +=
(UTF8String.fromString(key) -> convertField(parser, valueType, key))
case c: Characters if !c.isWhiteSpace =>
// Create a value tag field for it
kvPairs +=
// TODO: We don't support array value tags in maps yet.
(UTF8String.fromString(options.valueTag) -> convertTo(c.getData, valueType))
case _: EndElement | _: EndDocument =>
shouldStop = true
case _ => // do nothing
attributes: Array[Attribute],
keyType: DataType): MapData = {
val bounded = new ElementBoundedEventReader(parser)
try {
val kvPairs = ArrayBuffer.empty[(UTF8String, Any)]
attributes.foreach { attr =>
val key = convertXmlMapKey(
options.attributePrefix + attr.getName.getLocalPart, keyType)
kvPairs += (key -> convertTo(attr.getValue, valueType))
}
var shouldStop = false
while (!shouldStop) {
bounded.nextEvent match {
case e: StartElement =>
val rawName = StaxXmlParserUtils.getName(e.asStartElement.getName, options)
val key = convertXmlMapKey(rawName, keyType)
kvPairs += (key -> convertField(bounded, valueType, rawName))
case c: Characters if !c.isWhiteSpace =>
// Create a value tag field for it
kvPairs +=
// TODO: We don't support array value tags in maps yet.
(convertXmlMapKey(options.valueTag, keyType) -> convertTo(c.getData, valueType))
case _: EndElement | _: EndDocument =>
shouldStop = true
case _ => // do nothing
}
}
ArrayBasedMapData(kvPairs.toMap)
} finally {
bounded.drain()
}
ArrayBasedMapData(kvPairs.toMap)
}

private def convertConstrainedMap(
parser: XMLEventReader,
keyType: DataType,
valueType: DataType,
attributes: Array[Attribute]): MapData = {
val lastEntries =
mutable.LinkedHashMap.empty[UTF8String, (UTF8String, Option[Any])]
var badMapException: Option[Throwable] = None
def mapKey(raw: UTF8String): UTF8String = {
CharVarcharUtils.applyTextParseSemantics(raw, keyType)
}
def appendPair(rawKey: String, value: Option[Any]): Unit = {
try {
val rawKeyUtf8 = UTF8String.fromString(rawKey)
lastEntries.remove(rawKeyUtf8)
lastEntries.update(rawKeyUtf8, (mapKey(rawKeyUtf8), value))
} catch {
case NonFatal(e) => badMapException = badMapException.orElse(Some(e))
}
}
attributes.foreach { attr =>
val value = try {
Some(convertTo(attr.getValue, valueType))
} catch {
case e: SparkUpgradeException => throw e
case NonFatal(e) =>
badMapException = badMapException.orElse(Some(e))
None
}
appendPair(options.attributePrefix + attr.getName.getLocalPart, value)
}
var shouldStop = false
while (!shouldStop) {
parser.nextEvent match {
case e: StartElement =>
val rawKey = StaxXmlParserUtils.getName(e.asStartElement.getName, options)
val entryParser = new ElementBoundedEventReader(parser)
val value = {
try {
Some(convertField(entryParser, valueType, rawKey))
} catch {
case e: SparkUpgradeException => throw e
case DuplicateMapKeyUtils(e) => throw e
case NonFatal(e) =>
badMapException = badMapException.orElse(Some(e))
None
} finally {
entryParser.drain()
}
}
appendPair(rawKey, value)
case c: Characters if !c.isWhiteSpace =>
// Create a value tag field for it
// TODO: We don't support array value tags in maps yet.
val value = try {
Some(convertTo(c.getData, valueType))
} catch {
case e: SparkUpgradeException => throw e
case NonFatal(e) =>
badMapException = badMapException.orElse(Some(e))
None
}
appendPair(options.valueTag, value)
case _: EndElement | _: EndDocument =>
shouldStop = true
case _ => // do nothing
/**
* Length-check a CHAR/VARCHAR XML name used as a map key. Unlike value assignment,
* this does not pad or trim. STRING keys and the flag-off path still use
* [[CharVarcharUtils.applyTextParseSemantics]].
*/
private def convertXmlMapKey(rawName: String, keyType: DataType): UTF8String = {
val key = UTF8String.fromString(rawName)
if (charVarcharStandardSemantics) {
keyType match {
case c: CharType if key.numChars() != c.length =>
throw QueryExecutionErrors.unsupportedXmlCharVarcharMapKey(key, c)
case _: CharType => key
case v: VarcharType if key.numChars() > v.length =>
throw QueryExecutionErrors.unsupportedXmlCharVarcharMapKey(key, v)
case _: VarcharType => key
case _ => CharVarcharUtils.applyTextParseSemantics(key, keyType)
}
} else {
CharVarcharUtils.applyTextParseSemantics(key, keyType)
}
val mapData = DuplicateMapKeyUtils.buildConstrainedMap(
lastEntries, keyType, valueType)
badMapException.foreach(throw _)
mapData
}

/**
Expand Down Expand Up @@ -583,7 +551,6 @@ class StaxXmlParser(
row(i) = convertTo(v, schema(i).dataType)
} catch {
case e: SparkUpgradeException => throw e
case DuplicateMapKeyUtils(e) => throw e
case NonFatal(e) => firstError = firstError.orElse(Some(e))
}
}
Expand Down Expand Up @@ -722,7 +689,6 @@ class StaxXmlParser(
}
} catch {
case e: SparkUpgradeException => throw e
case DuplicateMapKeyUtils(e) => throw e
case NonFatal(e) =>
// TODO: we don't support partial results now
badRecordException = badRecordException.orElse(Some(e))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2739,6 +2739,20 @@ private[sql] object QueryExecutionErrors extends QueryErrorsBase with ExecutionE
)
}

/**
* Restricting CHAR/VARCHAR XML map keys is a parser limitation (SQLSTATE 0A000),
* but this is a SparkRuntimeException so XML parse modes can handle it: PERMISSIVE
* wraps a bad record and FAILFAST surfaces the failure.
*/
def unsupportedXmlCharVarcharMapKey(
key: UTF8String, dataType: DataType): SparkRuntimeException = {
new SparkRuntimeException(
errorClass = "UNSUPPORTED_XML_CHAR_VARCHAR_MAP_KEY",
messageParameters = Map(
"key" -> toSQLValue(key, StringType),
"dataType" -> toSQLType(dataType)))
}

def timestampAddOverflowError(micros: Long, amount: Long, unit: String): ArithmeticException = {
new SparkArithmeticException(
errorClass = "DATETIME_OVERFLOW",
Expand Down
Loading