Repository navigation
[SPARK-60108][SQL] Length-check JSON CHAR/VARCHAR map keys without pad or trim #59325
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
06ca409
ac8a4a7
d11523a
b9ee4f0
40ea574
6943ecc
ff58f3c
1425458
87bfa8c
4b2ab0d
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -1873,7 +1873,14 @@ case class JsonToStructs( | |
| options: Map[String, String], | ||
| child: Expression, | ||
| timeZoneId: Option[String] = None, | ||
| variantAllowDuplicateKeys: Boolean = SQLConf.get.getConf(SQLConf.VARIANT_ALLOW_DUPLICATE_KEYS)) | ||
| variantAllowDuplicateKeys: Boolean = SQLConf.get.getConf(SQLConf.VARIANT_ALLOW_DUPLICATE_KEYS), | ||
| // `spark.sql.charVarchar.standardSemantics.enabled` has PERSISTED binding, so it is captured | ||
| // here at analysis time (as the default arg, like `variantAllowDuplicateKeys`) and then | ||
| // threaded all the way into the parser via `charVarcharStandardSemanticsOverride`. This keeps | ||
| // a view's CHAR/VARCHAR map-key semantics tied to its creation-time flag rather than the | ||
| // caller's session setting. (`variantAllowDuplicateKeys` is only captured, not threaded: the | ||
| // parser still reads that one live from `SQLConf.get`.) | ||
| charVarcharStandardSemantics: Boolean = SQLConf.get.charVarcharStandardSemantics) | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Nit (P3): This default is evaluated wherever Take a no-SerDe Verification:
|
||
| extends UnaryExpression | ||
| with TimeZoneAwareExpression | ||
| with CodegenFallback | ||
|
|
@@ -1935,7 +1942,8 @@ case class JsonToStructs( | |
|
|
||
| @transient | ||
| private lazy val evaluator = new JsonToStructsEvaluator( | ||
| options, nullableSchema, nameOfCorruptRecord, timeZoneId, variantAllowDuplicateKeys) | ||
| options, nullableSchema, nameOfCorruptRecord, timeZoneId, variantAllowDuplicateKeys, | ||
| charVarcharStandardSemantics) | ||
| override def stateful: Boolean = true | ||
|
|
||
| override def nullSafeEval(json: Any): Any = { | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -51,7 +51,13 @@ class JacksonParser( | |
| schema: DataType, | ||
| val options: JSONOptions, | ||
| allowArrayAsStructs: Boolean, | ||
| filters: Seq[Filter] = Seq.empty) extends Logging { | ||
| filters: Seq[Filter] = Seq.empty, | ||
| // `spark.sql.charVarchar.standardSemantics.enabled` has PERSISTED binding, so for `from_json` | ||
| // it must be captured when the expression is analyzed (see `JsonToStructs`) rather than read | ||
| // live here, otherwise a view created under the flag would skip the CHAR/VARCHAR map-key check | ||
| // when queried from a session with the flag off. `None` falls back to the live conf, which is | ||
| // correct for the file-based JSON data source where there is no view to capture. | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Nit (P3): "there is no view to capture" does not hold for file scans. A persisted view can read a catalog JSON table, and under standard semantics that table keeps -- created with the flag on
CREATE TABLE t (m MAP<CHAR(3), INT>) USING json; -- data row: {"m":{"a":1}}
CREATE VIEW v AS SELECT m FROM t;
SELECT * FROM v;
-- queried with the flag on: m is NULL
-- queried with the flag off: m = {'a ' -> 1}In the flag-off session the parser installs no checker and accepts Verification:
|
||
| charVarcharStandardSemanticsOverride: Option[Boolean] = None) extends Logging { | ||
|
|
||
| import JacksonUtils._ | ||
| import com.fasterxml.jackson.core.JsonToken._ | ||
|
|
@@ -60,6 +66,13 @@ class JacksonParser( | |
| // to a value in a field for `InternalRow`. | ||
| private type ValueConverter = JsonParser => AnyRef | ||
|
|
||
| // CHAR/VARCHAR JSON map keys are length-checked without pad or trim under this flag. Lazy so | ||
| // its value is independent of field declaration order: `makeMapKeyChecker` reads it while the | ||
| // converters below are built, and an eager val read before its own initializer would silently | ||
| // see `false` and disable the check. | ||
| private lazy val charVarcharStandardSemantics = | ||
| charVarcharStandardSemanticsOverride.getOrElse(SQLConf.get.charVarcharStandardSemantics) | ||
|
|
||
| // `ValueConverter`s for the root schema for all fields in the schema | ||
| private val rootConverter = makeRootConverter(schema) | ||
|
|
||
|
|
@@ -191,8 +204,9 @@ class JacksonParser( | |
|
|
||
| private def makeMapRootConverter(mt: MapType): JsonParser => Iterable[InternalRow] = { | ||
| val fieldConverter = makeConverter(mt.valueType) | ||
| val keyChecker = makeMapKeyChecker(mt.keyType) | ||
| (parser: JsonParser) => parseJsonToken[Iterable[InternalRow]](parser, mt) { | ||
| case START_OBJECT => Some(InternalRow(convertMap(parser, fieldConverter))) | ||
| case START_OBJECT => Some(InternalRow(convertMap(parser, fieldConverter, keyChecker))) | ||
| } | ||
| } | ||
|
|
||
|
|
@@ -484,8 +498,9 @@ class JacksonParser( | |
|
|
||
| case mt: MapType => | ||
| val valueConverter = makeConverter(mt.valueType) | ||
| val keyChecker = makeMapKeyChecker(mt.keyType) | ||
| (parser: JsonParser) => parseJsonToken[MapData](parser, dataType) { | ||
| case START_OBJECT => convertMap(parser, valueConverter) | ||
| case START_OBJECT => convertMap(parser, valueConverter, keyChecker) | ||
| } | ||
|
|
||
| case udt: UserDefinedType[_] => | ||
|
|
@@ -618,38 +633,65 @@ class JacksonParser( | |
| } | ||
|
|
||
| /** | ||
| * Parse an object as a Map, preserving all fields. | ||
| * Parse an object as a Map. | ||
| * | ||
| * JSON object names used as CHAR/VARCHAR keys are length-checked without rewriting (see | ||
| * `makeMapKeyChecker`): 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, and duplicate names are kept, exactly as for STRING keys. This differs from XML, | ||
| * which pads keys and then applies `mapKeyDedupPolicy`. | ||
| * | ||
| * A rejected key is treated like a failed value, regardless of `enablePartialResults`: the | ||
| * cause is recorded and the parser steps past the value so the loop still consumes the map's | ||
| * END_OBJECT, leaving the key unpaired. The check must not surface its result while the parser | ||
| * is still on the FIELD_NAME, or an enclosing struct would misread the map's remaining entries | ||
| * as its own sibling fields (SPARK-60108). STRING keys have no checker and take the fast path. | ||
| */ | ||
| private def convertMap( | ||
| parser: JsonParser, | ||
| fieldConverter: ValueConverter): MapData = { | ||
| fieldConverter: ValueConverter, | ||
| keyChecker: Option[UTF8String => Option[Throwable]]): 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) | ||
| try { | ||
| values += fieldConverter.apply(parser) | ||
| } catch { | ||
| case err: PartialValueException if enablePartialResults => | ||
| badRecordException = badRecordException.orElse(Some(err.cause)) | ||
| values += err.partialResult | ||
| case NonFatal(e) if enablePartialResults => | ||
| val key = UTF8String.fromString(parser.currentName) | ||
| keys += key | ||
| // Avoid allocating a closure per entry on the common no-checker path. | ||
| val keyRejection = keyChecker match { | ||
| case Some(check) => check(key) | ||
| case None => None | ||
| } | ||
| keyRejection match { | ||
| case Some(e) => | ||
| badRecordException = badRecordException.orElse(Some(e)) | ||
| parser.nextToken() // step from FIELD_NAME onto the value | ||
| parser.skipChildren() | ||
| case None => | ||
| try { | ||
| values += fieldConverter.apply(parser) | ||
| } catch { | ||
| case err: PartialValueException if enablePartialResults => | ||
| badRecordException = badRecordException.orElse(Some(err.cause)) | ||
| values += err.partialResult | ||
| case NonFatal(e) if enablePartialResults => | ||
| badRecordException = badRecordException.orElse(Some(e)) | ||
| parser.skipChildren() | ||
| } | ||
| } | ||
| } | ||
|
|
||
| // Value conversion can fail after the key is recorded. Do not build MapData from | ||
| // unpaired buffers: ArrayBasedMapData would throw a cardinality error and hide | ||
| // the original conversion failure (for example EXCEED_LIMIT_LENGTH). | ||
| // A rejected key or a failed value leaves the key recorded without a value. Rethrow the | ||
| // real cause instead of building an unbalanced ArrayBasedMapData, whose cardinality | ||
| // `require` would otherwise mask it. The row then becomes a bad record (null in PERMISSIVE, | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Nit (P3): This says the row then becomes a bad record, Verification:
|
||
| // surfaced in FAILFAST). Balanced partial results (a trimmed value) still flow through | ||
| // PartialMapDataResultException below. | ||
| if (keys.length != values.length) { | ||
| throw badRecordException.get | ||
| } | ||
|
|
||
| // Preserve every parsed JSON key/value pair, including exact duplicate names. | ||
| // ArrayBasedMapData is used directly to retain this historical behavior. | ||
| // The JSON map keeps every parsed pair, including exact duplicate names. | ||
| val mapData = ArrayBasedMapData(keys.toArray, values.toArray) | ||
|
|
||
| if (badRecordException.isEmpty) { | ||
|
|
@@ -659,6 +701,44 @@ class JacksonParser( | |
| } | ||
| } | ||
|
|
||
| /** | ||
| * Builds the length check applied to this map's CHAR/VARCHAR keys, or `None` when the keys | ||
| * need no check (STRING keys, or standard semantics are off). Built once per map converter so | ||
| * the per-entry loop in `convertMap` does not re-inspect the key type. The check never pads or | ||
| * trims; it returns the rejection cause (with the actual key type, collation included, so the | ||
| * message is accurate) rather than throwing, so `convertMap` keeps full control of the parser | ||
| * position when a key is rejected. | ||
| */ | ||
| private def makeMapKeyChecker(keyType: DataType): Option[UTF8String => Option[Throwable]] = { | ||
| if (!charVarcharStandardSemantics) { | ||
| None | ||
| } else { | ||
| keyType match { | ||
| case c: CharType => Some(checkCharJsonMapKey(_, c)) | ||
| case v: VarcharType => Some(checkVarcharJsonMapKey(_, v)) | ||
| case _ => None | ||
| } | ||
| } | ||
| } | ||
|
|
||
| // A CHAR(n) JSON object name is never padded: it must already be exactly n characters. | ||
| private def checkCharJsonMapKey(key: UTF8String, keyType: CharType): Option[Throwable] = { | ||
| if (key.numChars() != keyType.length) { | ||
| Some(QueryExecutionErrors.unsupportedJsonCharVarcharMapKey(key, keyType)) | ||
| } else { | ||
| None | ||
| } | ||
| } | ||
|
|
||
| // A VARCHAR(n) JSON object name is never trimmed: it must already be at most n characters. | ||
| private def checkVarcharJsonMapKey(key: UTF8String, keyType: VarcharType): Option[Throwable] = { | ||
| if (key.numChars() > keyType.length) { | ||
| Some(QueryExecutionErrors.unsupportedJsonCharVarcharMapKey(key, keyType)) | ||
| } else { | ||
| None | ||
| } | ||
| } | ||
|
|
||
| /** | ||
| * Parse an object as a Array. | ||
| */ | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Non-blocking (P2): The last sentence here says that reads whose CHAR/VARCHAR keys are handled read-side, "for example catalog tables", are unchanged: keys are padded,
spark.sql.mapKeyDedupPolicyis applied, and long keys raiseEXCEED_LIMIT_LENGTH. That is not what happens under standard semantics. The sentence comes from my earlier comment, which had this wrong.With
spark.sql.charVarchar.standardSemantics.enabled=true,charVarcharFirstClassTypesis true, soreplaceCharVarcharWithStringkeepsCharType. For a catalog JSON table,SessionCatalog.getTableMetadataandLogicalRelation.applykeepMAP<CHAR(3), INT>,DataSourceStrategy.readDataSourceTablepasses it as the reader schema, andJsonFileFormatbuildsJacksonParserover it.makeMapKeyCheckerthen installs the new check:Neither row gets the documented
'a 'orEXCEED_LIMIT_LENGTH. Keys that pass the parser are then rebuilt by the read-sideMapFromArrays, for catalog tables and user-specified reader schemas alike. So{"m":{"abc":1,"abc":2}}still raisesDUPLICATED_MAP_KEYunder the default policy, unlikefrom_json. Could this sentence, and the matching PR-description paragraph, say that JSON scans with CHAR/VARCHAR map keys run the same parser check and still applymapKeyDedupPolicy? If catalog reads should skip the parser check instead, that would need a code change.See Shared repair plan 1 in the review body.