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

Spark中groupBy结合collect_list(struct)报错及JVM崩溃问题求助

根本原因

你遇到的是Spark 2.3.x版本Java Bean Encoder的已知缺陷:该版本的Java Bean编码器在处理嵌套复杂类型(如List<自定义POJO>对应Spark的ArrayType(StructType(...)))时,无法正确对齐嵌套字段的字节偏移量,读取数据时会出现字节错位:

  • 当嵌套字段包含8字节的Double类型时,错位后会读取到无效的数组长度值,触发valueArraySize必须为正或超过Integer.MAX_VALUE的报错,严重时直接导致JVM崩溃
  • 当嵌套字段均为4字节的Integer类型时,字节错位不会触发直接越界,但会出现字段取值错位,导致输出结果不符合预期

可行解决方案

方案1:手动映射Row对象(最稳妥,无需改版本/依赖)

避免直接用Encoders.bean(Top.class)做类型转换,先收集Row类型结果,手动遍历映射为自定义类,代码示例:

List<Row> rowList = grouped.collectAsList();
List<Top> output = new ArrayList<>();
for (Row row : rowList) {
    Top top = new Top();
    top.setUser(row.getString(row.fieldIndex("user")));
    List<Command> commands = new ArrayList<>();
    List<Row> commandRows = row.getList(row.fieldIndex("commands"));
    for (Row commandRow : commandRows) {
        Command cmd = new Command();
        cmd.setItem(commandRow.getString(commandRow.fieldIndex("item")));
        cmd.setAmount(commandRow.getDouble(commandRow.fieldIndex("amount")));
        commands.add(cmd);
    }
    top.setCommands(commands);
    output.add(top);
}

方案2:升级Spark版本

该嵌套类型序列化bug在Spark 2.4.0及以上版本已被官方修复,若业务允许升级Spark,直接将版本升级到2.4.0或更高版本即可解决问题。

方案3:通过JSON中转适配

如果既不能升级Spark,也不想手动写Row映射,可以先将嵌套struct列序列化为JSON字符串,再用Jackson等JSON工具反序列化为对应POJO:

// 聚合后先转JSON
Dataset<Row> grouped = spark.createDataset(orders, Encoders.bean(Order.class))
        .groupBy("user")
        .agg(collect_list(struct( "item", "amount" )).as("commands"))
        .withColumn("commands", to_json(col("commands")));

// 收集后反序列化
ObjectMapper objectMapper = new ObjectMapper();
List<Top> output = new ArrayList<>();
for (Row row : grouped.collectAsList()) {
    Top top = new Top();
    top.setUser(row.getString(0));
    String commandsJson = row.getString(1);
    List<Command> commands = objectMapper.readValue(commandsJson, new TypeReference<List<Command>>() {});
    top.setCommands(commands);
    output.add(top);
}

内容的提问来源于stack exchange,提问作者Juh_

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 10:39:00