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_
相关产品推荐
相关产品推荐

