You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.14 17:14:34