srielau commented on code in PR #58545:
URL: https://github.com/apache/spark/pull/58545#discussion_r3994128737
##########
sql/core/src/test/scala/org/apache/spark/sql/CharVarcharTestSuite.scala:
##########
@@ -2226,6 +2251,138 @@ class BasicCharVarcharTestSuite extends
SharedSparkSession {
}
}
}
+
+ test("SPARK-59274: from_json/csv/xml honor CHAR/VARCHAR under
standardSemantics") {
+ withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true") {
+ val jsonChar = sql("""SELECT from_json('{"a": "str"}', 'a CHAR(5)')""")
+ val jsonCharType = jsonChar.schema.head.dataType.asInstanceOf[StructType]
+ assert(jsonCharType.head.dataType === CharType(5))
+ checkAnswer(jsonChar, Row(Row("str ")))
+
+ val jsonVarchar = sql("""SELECT from_json('{"a": "ab"}', 'a
VARCHAR(5)')""")
+ val jsonVarcharType =
jsonVarchar.schema.head.dataType.asInstanceOf[StructType]
+ assert(jsonVarcharType.head.dataType === VarcharType(5))
+ checkAnswer(jsonVarchar, Row(Row("ab")))
+
+ // Default PERMISSIVE mode turns length failures into a null record.
+ Seq("CHAR(5)", "VARCHAR(5)").foreach { dataType =>
+ checkAnswer(
+ sql(s"""SELECT from_json('{"a": "abcdef"}', 'a $dataType')"""),
+ Row(Row(null)))
+ assertParseExceedLimit(
+ s"""SELECT from_json(
+ | '{"a": "abcdef"}',
+ | 'a $dataType',
+ | map('mode', 'FAILFAST'))""".stripMargin)
+ }
+
+ checkAnswer(
+ sql("""SELECT from_json('{"ab": 1}', 'MAP<CHAR(4), INT>')"""),
+ Row(Map("ab " -> 1)))
+
+ checkAnswer(sql("SELECT from_csv('str', 'a CHAR(5)')"), Row(Row("str
")))
+ Seq("CHAR(5)", "VARCHAR(5)").foreach { dataType =>
+ checkAnswer(sql(s"SELECT from_csv('abcdef', 'a $dataType')"),
Row(Row(null)))
+ assertParseExceedLimit(
+ s"SELECT from_csv('abcdef', 'a $dataType', map('mode', 'FAILFAST'))")
+ }
+
+ checkAnswer(
+ sql("SELECT from_xml('<ROW><a>str</a></ROW>', 'a CHAR(5)')"),
+ Row(Row("str ")))
+ checkAnswer(
+ sql(
+ """SELECT from_xml(
+ | '<ROW><a></a></ROW>',
+ | 'a CHAR(5)',
+ | map('nullValue', 'NULL'))""".stripMargin),
+ Row(Row(" ")))
+ Seq("CHAR(5)", "VARCHAR(5)").foreach { dataType =>
+ checkAnswer(
+ sql(s"SELECT from_xml('<ROW><a>abcdef</a></ROW>', 'a $dataType')"),
+ Row(Row(null)))
+ assertParseExceedLimit(
+ s"SELECT from_xml('<ROW><a>abcdef</a></ROW>', 'a $dataType', " +
+ "map('mode', 'FAILFAST'))")
+ }
+ checkAnswer(
+ sql("SELECT from_xml('<ROW><m><ab>1</ab></m></ROW>', 'm MAP<CHAR(4),
INT>')"),
+ Row(Row(Map("ab " -> 1))))
+
+ checkAnswer(
+ sql("""SELECT schema_of_json(CAST('{"a":1}' AS VARCHAR(20)))"""),
+ Row("STRUCT<a: BIGINT>"))
+ checkAnswer(
+ sql("SELECT schema_of_csv(CAST('1,abc' AS VARCHAR(20)))"),
+ Row("STRUCT<_c0: INT, _c1: STRING>"))
+ checkAnswer(
+ sql("SELECT schema_of_xml(CAST('<ROW><a>1</a></ROW>' AS
VARCHAR(40)))"),
+ Row("STRUCT<a: BIGINT>"))
+ }
+ }
+
+ test("SPARK-59274: normalized CHAR map key collisions honor the dedup
policy") {
+ withSQLConf(
+ SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true",
Review Comment:
Fixed in 0076ea48764. Added VARCHAR normalized-collision and overflow
coverage, a normalized duplicate in FAILFAST mode, and an XML file-source read
using a constrained user schema.
##########
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/xml/StaxXmlParser.scala:
##########
@@ -374,31 +378,52 @@ class StaxXmlParser(
*/
private def convertMap(
parser: XMLEventReader,
+ keyType: DataType,
valueType: DataType,
attributes: Array[Attribute]): MapData = {
val kvPairs = ArrayBuffer.empty[(UTF8String, Any)]
+ var mapKeyException: Option[Throwable] = None
+ def mapKey(raw: String): UTF8String = {
+ CharVarcharUtils.applyTextParseSemantics(UTF8String.fromString(raw),
keyType)
+ }
+ def appendPair(rawKey: String, value: Any): Unit = {
+ try {
+ kvPairs += (mapKey(rawKey) -> value)
+ } catch {
+ case NonFatal(e) => mapKeyException = mapKeyException.orElse(Some(e))
+ }
+ }
attributes.foreach { attr =>
- kvPairs += (UTF8String.fromString(options.attributePrefix +
attr.getName.getLocalPart)
- -> convertTo(attr.getValue, valueType))
+ val value = convertTo(attr.getValue, valueType)
+ appendPair(options.attributePrefix + attr.getName.getLocalPart, value)
}
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))
+ val rawKey = StaxXmlParserUtils.getName(e.asStartElement.getName,
options)
+ val value = convertField(parser, valueType, rawKey)
+ appendPair(rawKey, value)
case c: Characters if !c.isWhiteSpace =>
// Create a value tag field for it
- kvPairs +=
// TODO: We don't support an array value tags in map yet.
- (UTF8String.fromString(options.valueTag) -> convertTo(c.getData,
valueType))
+ val value = convertTo(c.getData, valueType)
+ appendPair(options.valueTag, value)
case _: EndElement | _: EndDocument =>
shouldStop = true
case _ => // do nothing
}
}
- ArrayBasedMapData(kvPairs.toMap)
+ mapKeyException.foreach(throw _)
Review Comment:
Fixed in 0076ea48764. Constrained XML maps now build first so duplicate
policy takes precedence, then surface any remembered non-duplicate key error.
The regression covers an overlength key followed by a normalized collision.
##########
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/json/JacksonParser.scala:
##########
@@ -617,16 +629,20 @@ class JacksonParser(
*/
private def convertMap(
parser: JsonParser,
- fieldConverter: ValueConverter): MapData = {
+ fieldConverter: ValueConverter,
+ keyType: DataType,
+ valueType: DataType): MapData = {
val keys = ArrayBuffer.empty[UTF8String]
val values = ArrayBuffer.empty[Any]
var badRecordException: Option[Throwable] = None
while (nextUntil(parser, JsonToken.END_OBJECT)) {
- keys += UTF8String.fromString(parser.currentName)
+ keys += CharVarcharUtils.applyTextParseSemantics(
Review Comment:
Fixed in 0076ea48764. Normalized-key tracking is now separate from retained
key/value pairs, so the malformed first value no longer suppresses EXCEPTION
while partial-map materialization remains atomic.
##########
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/util/CharVarcharUtils.scala:
##########
@@ -156,6 +157,25 @@ object CharVarcharUtils extends Logging with
SparkCharVarcharUtils {
StructType(fields)
}
+ /**
Review Comment:
Fixed in 0076ea48764. The comment now describes assignment semantics
precisely: CHAR padding, permitted excess trailing-space trimming, and
EXCEED_LIMIT_LENGTH only for remaining non-space overflow.
##########
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/json/JacksonParser.scala:
##########
@@ -617,33 +629,58 @@ class JacksonParser(
*/
private def convertMap(
parser: JsonParser,
- fieldConverter: ValueConverter): MapData = {
+ fieldConverter: ValueConverter,
+ keyType: DataType,
+ valueType: DataType): MapData = {
val keys = ArrayBuffer.empty[UTF8String]
val values = ArrayBuffer.empty[Any]
- var badRecordException: Option[Throwable] = None
+ var partialResultException: Option[Throwable] = None
+ var badMapException: Option[Throwable] = None
while (nextUntil(parser, JsonToken.END_OBJECT)) {
- keys += UTF8String.fromString(parser.currentName)
- try {
- values += fieldConverter.apply(parser)
+ val rawKey = UTF8String.fromString(parser.currentName)
+ val value = try {
+ Some(fieldConverter.apply(parser))
} catch {
case err: PartialValueException if enablePartialResults =>
- badRecordException = badRecordException.orElse(Some(err.cause))
- values += err.partialResult
+ partialResultException =
partialResultException.orElse(Some(err.cause))
+ Some(err.partialResult)
+ case DuplicateMapKeyException(e) => throw e
case NonFatal(e) if enablePartialResults =>
- badRecordException = badRecordException.orElse(Some(e))
+ badMapException = badMapException.orElse(Some(e))
parser.skipChildren()
+ None
+ }
+ value.foreach { parsedValue =>
Review Comment:
Fixed in 0076ea48764. JSON now records every successfully normalized
constrained key after consuming its value, independently of whether that value
is retained in the partial map. It applies the duplicate policy to that full
key stream before building the retained map. The exact malformed-value
collision is covered under EXCEPTION and LAST_WIN.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]