[
https://issues.apache.org/jira/browse/SPARK-51070?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated SPARK-51070:
-----------------------------------
Labels: pull-request-available (was: )
> Expression encoding fails on Array/Seq/Map with mutable.Set
> -----------------------------------------------------------
>
> Key: SPARK-51070
> URL: https://issues.apache.org/jira/browse/SPARK-51070
> Project: Spark
> Issue Type: Bug
> Components: SQL
> Affects Versions: 4.0.0, 3.4.4, 3.5.4
> Reporter: Chaerim Yeo
> Priority: Major
> Labels: pull-request-available
>
> Although this is a rare situation, the following error occurs when trying to
> create a {{Dataset}} from {{{}Seq[mutable.Set]{}}}:
> {code:scala}
> scala> val set: collection.Set[Int] = mutable.Set(1, 2)
> set: scala.collection.Set[Int] = Set(1, 2)
> scala> Seq(Seq(set)).toDS()
> org.apache.spark.SparkRuntimeException: Error while encoding:
> java.lang.RuntimeException: scala.collection.mutable.HashSet is not a valid
> external type for schema of array<int>
> mapobjects(lambdavariable(MapObject, ObjectType(class java.lang.Object),
> true, -1), mapobjects(lambdavariable(MapObject, ObjectType(class
> java.lang.Object), true, -2),
> assertnotnull(validateexternaltype(lambdavariable(MapObject, ObjectType(class
> java.lang.Object), true, -2), IntegerType, IntegerType)),
> validateexternaltype(lambdavariable(MapObject, ObjectType(class
> java.lang.Object), true, -1), ArrayType(IntegerType,false),
> ObjectType(interface scala.collection.Set)).toSeq, None), input[0,
> scala.collection.Seq, true], None) AS value#16.
> at
> org.apache.spark.sql.errors.QueryExecutionErrors$.expressionEncodingError(QueryExecutionErrors.scala:1561)
> at
> org.apache.spark.sql.catalyst.encoders.ExpressionEncoder$Serializer.apply(ExpressionEncoder.scala:210)
> at
> org.apache.spark.sql.SparkSession.$anonfun$createDataset$1(SparkSession.scala:483)
> at scala.collection.immutable.List.map(List.scala:293)
> at org.apache.spark.sql.SparkSession.createDataset(SparkSession.scala:483)
> at org.apache.spark.sql.SQLContext.createDataset(SQLContext.scala:354)
> at
> org.apache.spark.sql.SQLImplicits.localSeqToDatasetHolder(SQLImplicits.scala:244)
> ... 51 elided
> Caused by: java.lang.RuntimeException: scala.collection.mutable.HashSet is
> not a valid external type for schema of array<int>
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.Invoke_0$(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.MapObjects_0$(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.MapObjects_1$(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.apply(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.encoders.ExpressionEncoder$Serializer.apply(ExpressionEncoder.scala:207)
> ... 56 more
> {code}
> According to the codegen result, it seems that
> `Scala.collection.immutable.Set` was used for validation.
> {code:java}
> /* 092 */ private scala.collection.Seq Invoke_0(InternalRow i) {
> /* 093 */ scala.collection.Set value_4 = null;
> /* 094 */ if (!isNull_MapObject_lambda_variable_4) {
> /* 095 */ if (value_MapObject_lambda_variable_4.getClass().isArray() ||
> value_MapObject_lambda_variable_4 instanceof scala.collection.Seq ||
> value_MapObject_lambda_variable_4 instanceof scala.collection.immutable.Set
> || value_MapObject_lambda_variable_4 instanceof java.util.List) {
> /* 096 */ value_4 = (scala.collection.Set)
> value_MapObject_lambda_variable_4;
> /* 097 */ } else {
> /* 098 */ throw new
> RuntimeException(value_MapObject_lambda_variable_4.getClass().getName() +
> ((java.lang.String) references[0] /* errMsg */));
> /* 099 */ }
> /* 100 */ }
> {code}
> This issue is caused by {{ValidateExternalType}} and occurs on creating
> {{Dataset}} from {{{}Array{}}}/{{{}Seq{}}}/{{{}Map{}}} with
> {{{}mutable.Set{}}}:
> * When constructing {{{}AgnosticEncoder{}}}, {{scala.collection.Set}} is
> converted into {{{}IterableEncoder{}}}.
> * When creating a serializer for {{{}ExpressionEncoder{}}}, the following
> elements are wrapped with {{ValidateExternalType}} by
> {{{}SerializerBuildHelper.validateAndSerializeElement{}}}:
> ** Element type of {{ArrayEncoder}}
> ** Element type of {{IterableEncoder}}
> ** Key type of {{MapEncoder}}
> ** Value type of {{MapEncoder}}
> * {{ValidateExternalType}} attempts to validate using {{{}Set{}}}, which is
> equivalent to {{scala.collection.immutable.Set}} (regardness of Scala 2.12 /
> 2.13).
> {code:scala}
> scala> :paste
> // Entering paste mode (ctrl-D to finish)
> import org.apache.spark.sql.catalyst.encoders.ExpressionEncoder
> import org.apache.spark.sql.catalyst.expressions.codegen._
> import org.apache.spark.sql.catalyst.expressions.{GenericInternalRow,
> UnsafeRow}
> import scala.collection.mutable
> import scala.reflect.runtime.universe.TypeTag
> import scala.util.control.Exception.allCatch
> def test[T: TypeTag](value: T): UnsafeRow = {
> val encoder = ExpressionEncoder[T]
> val projection = GenerateUnsafeProjection.generate(encoder.serializer,
> subexpressionEliminationEnabled = true)
> val inputRow = new GenericInternalRow(1)
> inputRow(0) = value
> projection(inputRow)
> }
> val set: collection.Set[Int] = mutable.Set(1, 2)
> // Exiting paste mode, now interpreting.
> import org.apache.spark.sql.catalyst.encoders.ExpressionEncoder
> import org.apache.spark.sql.catalyst.expressions.codegen._
> import org.apache.spark.sql.catalyst.expressions.{GenericInternalRow,
> UnsafeRow}
> import scala.collection.mutable
> import scala.reflect.runtime.universe.TypeTag
> import scala.util.control.Exception.allCatch
> test: [T](value: T)(implicit evidence$1:
> reflect.runtime.universe.TypeTag[T])org.apache.spark.sql.catalyst.expressions.UnsafeRow
> set: scala.collection.Set[Int] = Set(1, 2)
> scala> test(set)
> res0: org.apache.spark.sql.catalyst.expressions.UnsafeRow =
> [0,1000000018,2,0,200000001]
> scala> test(Array(set))
> java.lang.RuntimeException: scala.collection.mutable.HashSet is not a valid
> external type for schema of array<int>
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.Invoke_0$(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.MapObjects_0$(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.MapObjects_1$(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.apply(Unknown
> Source)
> at test(<pastie>:43)
> ... 51 elided
> scala> test(Seq(set))
> java.lang.RuntimeException: scala.collection.mutable.HashSet is not a valid
> external type for schema of array<int>
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.Invoke_0$(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.MapObjects_0$(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.MapObjects_1$(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.apply(Unknown
> Source)
> at test(<pastie>:43)
> ... 51 elided
> scala> test(Map(set -> 1))
> java.lang.RuntimeException: scala.collection.mutable.HashSet is not a valid
> external type for schema of array<int>
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.Invoke_0$(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.MapObjects_0$(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.ExternalMapToCatalyst_0$(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.apply(Unknown
> Source)
> at test(<pastie>:43)
> ... 51 elided
> scala> test(Map(1 -> set))
> java.lang.RuntimeException: scala.collection.mutable.HashSet is not a valid
> external type for schema of array<int>
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.Invoke_0$(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.MapObjects_0$(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.ExternalMapToCatalyst_0$(Unknown
> Source)
> at
> org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.apply(Unknown
> Source)
> at test(<pastie>:43)
> ... 51 elided
> {code}
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]