cloud-fan commented on code in PR #58545:
URL: https://github.com/apache/spark/pull/58545#discussion_r4037335814


##########
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:
   Confirmed: failed-value keys now participate in normalized duplicate 
validation without being materialized in the returned partial map.
   
   <!-- SPARK_DEV_REVIEW_REPLY 
{"feedback_id":"inline:3991524810","thread_id":"inline:3991524810","verdict_sha256":"1d67a83cffd94f727711a4c377b782e30aa3cf967fd5d486998fc278c510ae37"}
 -->



##########
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:
   Confirmed: the suite now covers VARCHAR key normalization and overflow, 
FAILFAST duplicate propagation, and the optimized XML file-source path.
   
   <!-- SPARK_DEV_REVIEW_REPLY 
{"feedback_id":"inline:3991524821","thread_id":"inline:3991524821","verdict_sha256":"1d67a83cffd94f727711a4c377b782e30aa3cf967fd5d486998fc278c510ae37"}
 -->



##########
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:
   Confirmed: the helper comment now distinguishes permitted trailing-space 
trimming from non-space overflow and describes CHAR padding accurately.
   
   <!-- SPARK_DEV_REVIEW_REPLY 
{"feedback_id":"inline:3991524825","thread_id":"inline:3991524825","verdict_sha256":"1d67a83cffd94f727711a4c377b782e30aa3cf967fd5d486998fc278c510ae37"}
 -->



##########
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:
   Confirmed: constrained XML maps now apply duplicate policy before surfacing 
a remembered non-duplicate key error.
   
   <!-- SPARK_DEV_REVIEW_REPLY 
{"feedback_id":"inline:3991524830","thread_id":"inline:3991524830","verdict_sha256":"1d67a83cffd94f727711a4c377b782e30aa3cf967fd5d486998fc278c510ae37"}
 -->



##########
sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/xml/StaxXmlParser.scala:
##########
@@ -374,31 +378,54 @@ 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)
+    keyType match {
+      case _: CharType | _: VarcharType =>
+        val mapBuilder = new ArrayBasedMapBuilder(keyType, valueType)
+        kvPairs.foreach { case (key, value) => mapBuilder.put(key, value) }
+        val mapData = mapBuilder.build()
+        mapKeyException.foreach(throw _)
+        mapData
+      case _ =>

Review Comment:
   I still see non-CHAR/VARCHAR XML keys routed through 
`collapseOrdinaryStringKeys` and Scala `Map.toMap`. For UTF8_LCASE, 
binary-distinct `a` and `A` therefore remain separate and bypass 
`mapKeyDedupPolicy`. Please route non-binary-collated STRING keys through the 
type-aware builder and update the current expectation.
   
   <!-- SPARK_DEV_REVIEW_REPLY 
{"feedback_id":"inline:4008384565","thread_id":"inline:4008384565","verdict_sha256":"1d67a83cffd94f727711a4c377b782e30aa3cf967fd5d486998fc278c510ae37"}
 -->



-- 
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]

Reply via email to