Spark 4 Java API调用try_variant_get时触发ClassCastException异常的问题求助
Spark 4 Java API调用try_variant_get时触发ClassCastException异常的问题求助
我最近在测试Spark 4中的try_variant_get方法处理Variant类型数据的能力,先通过SQL语句验证了功能是完全正常的,但转换成Java API实现时却抛出了ClassCastException,实在找不到问题所在,想请教大家怎么解决这个错误。
一、SQL语句验证(功能正常)
首先我创建了测试表并插入测试数据,SQL执行全程没有报错:
CREATE TABLE family ( id INT, data VARIANT ); INSERT INTO family (id, data) VALUES (1, PARSE_JSON('{"name":"Alice","age":30}')), (2, PARSE_JSON('[1,2,3,4,5]')), (3, PARSE_JSON('42'));
接着使用try_variant_get编写查询语句,成功返回了预期的输出结果:
SELECT id, try_variant_get(data, '$.name', 'STRING') AS name, try_variant_get(data, '$.age', 'INT') AS age FROM family WHERE try_variant_get(data, '$.name', 'STRING') IS NOT NULL;
二、转换为Java API后抛出异常
我把上述SQL逻辑转换成Spark Java API代码,但运行时直接触发了ClassCastException:
SparkSession spark = SparkSession.builder().master("local[*]").appName("VariantExample").getOrCreate(); StructType schema = new StructType() .add("id", DataTypes.IntegerType) .add("data", DataTypes.VariantType); Dataset<Row> df = spark.createDataFrame( Arrays.asList( RowFactory.create(1, "{\"name\":\"Alice\",\"age\":30}"), RowFactory.create(2, "[1,2,3,4,5]"), RowFactory.create(3, "42") ), schema ); Dataset<Row> df_sel = df.select( col("id"), try_variant_get(col("data"), "$.name", "String").alias("name"), try_variant_get(col("data"), "$.age", "Integer").alias("age") ).where("name IS NOT NULL"); df_sel.printSchema(); df_sel.show();
异常信息:
root |-- id: integer (nullable = true) |-- name: string (nullable = true) |-- age: integer (nullable = true) Exception in thread "main" java.lang.ClassCastException: class java.lang.String cannot be cast to class org.apache.spark.unsafe.types.VariantVal (java.lang.String is in module java.base of loader 'bootstrap'; org.apache.spark.unsafe.types.VariantVal is in unnamed module of loader 'app') at org.apache.spark.sql.catalyst.expressions.variant.VariantGet.nullSafeEval(variantExpressions.scala:282) at org.apache.spark.sql.catalyst.expressions.BinaryExpression.eval(Expression.scala:692) at org.apache.spark.sql.catalyst.expressions.Alias.eval(namedExpressions.scala:159) at org.apache.spark.sql.catalyst.expressions.InterpretedMutableProjection.apply(InterpretedMutableProjection.scala:89) at org.apache.spark.sql.catalyst.optimizer.ConvertToLocalRelation$$anonfun$apply$48.$anonfun$applyOrElse$83(Optimizer.scala:2231) at scala.collection.immutable.List.map(List.scala:247) at scala.collection.immutable.List.map(List.scala:79).....
我怀疑是try_variant_get方法的参数问题,但实在找不到具体哪里错了,希望大家能帮我看看怎么修复这个错误。
备注:内容来源于stack exchange,提问作者Joseph Hwang
相关产品推荐
相关产品推荐

