Spark 1.6(CDH5.9.2)中groupBy+collect_set致结构成员名小写问题
解决Spark 1.6中Case Class聚合后结构成员名小写的问题
我之前在Spark 1.x版本里也碰到过一模一样的字段名大小写坑,结合你用的CDH 5.9.2 + Spark 1.6环境,给你拆解下问题原因和可行的解决办法:
问题根源
Spark 1.6通过反射推断Scala Case Class的Schema时,会默认把Case Class的字段名转成小写。当你用collect_set聚合follower字段时,聚合后的集合元素会沿用这个被修改过的Schema,自然就出现了结构成员名变成小写的异常。
解决方案
方案1:手动指定Schema(推荐,无需改动现有Case Class)
既然反射推断会乱改大小写,那我们直接手动定义StructType来强制锁定字段名,确保Schema完全符合预期:
import org.apache.spark.sql.types.{StructType, StructField, StringType, LongType} import org.apache.spark.sql.functions._ // 手动定义UuidWrapper对应的Schema,严格保留驼峰字段名 val uuidWrapperSchema = StructType( Seq( StructField("id", StringType, nullable = true), StructField("lastSeenDate", LongType, nullable = true), StructField("firstSeenDate", LongType, nullable = true) ) ) // 先修复r1中leadR和follower字段的Schema val r1Fixed = r1.select( col("cid"), col("leadR").cast(uuidWrapperSchema), col("follower").cast(uuidWrapperSchema) ) // 再执行聚合操作 val r2 = r1Fixed.groupBy('cid, 'leadR).agg(collect_set('follower) as "followRz") // 验证结果 println("r2:") r2.show(20,false) r2.printSchema()
方案2:改用JavaBean替代Scala Case Class
Spark对JavaBean的Schema推断是严格遵循JavaBean规范(通过getter/setter方法映射字段名),不会自动转换大小写。你可以创建一个JavaBean类:
import java.io.Serializable; public class UuidWrapperBean implements Serializable { private String id; private Long lastSeenDate; private Long firstSeenDate; // 必须提供无参构造函数 public UuidWrapperBean() {} // Getter和Setter方法 public String getId() { return id; } public void setId(String id) { this.id = id; } public Long getLastSeenDate() { return lastSeenDate; } public void setLastSeenDate(Long lastSeenDate) { this.lastSeenDate = lastSeenDate; } public Long getFirstSeenDate() { return firstSeenDate; } public void setFirstSeenDate(Long firstSeenDate) { this.firstSeenDate = firstSeenDate; } }
然后在Scala代码中用这个JavaBean来构建r1的结构,后续聚合操作就能完美保留驼峰字段名了。
方案3:升级Spark版本(可选)
如果你的环境允许升级,Spark 2.x及以上版本已经修复了这个问题——默认会保留Scala Case Class的原字段名,不会自动转小写。不过考虑到你用的是绑定Spark 1.6的CDH 5.9.2,这个方案可能落地难度较高。
内容的提问来源于stack exchange,提问作者codeaperature
相关产品推荐
相关产品推荐

