Spark groupByKey无法识别已知字段lambda调用,为何需用getAs?
为什么Spark已知Schema却不能用
x.startId访问字段,必须用getAs? 核心原因:静态编译检查 vs Spark运行时Schema
你遇到的问题本质是Scala静态类型检查和Spark运行时元数据之间的差异:
joinedDS本质是Dataset[Row],属于弱类型数据集
不管Spark的Schema里定义了多少字段,Row类本身并没有startId、address这类成员属性。Scala是静态类型语言,编译阶段会直接检查x.startId是否合法——因为Row没有这个方法/属性,所以直接抛出编译错误,根本轮不到Spark去用Schema做推断。Spark的Schema是运行时元数据,编译期不可见
Spark的Schema信息是存储在Dataset的元数据里,只有在程序运行时才能访问。而Scala编译器在编译你的代码时,完全不知道Row里包含哪些字段,自然无法把x.startId转换成x.getAs[String]("startId")这种合法调用。getAs是Row类的合法方法Row.getAs[String]("startId")是Scala编译器能识别的合法调用:它调用了Row类自带的getAs方法,通过字段名获取对应类型的值。此时字段是否存在、类型是否匹配,才会在运行时由Spark根据Schema去校验。
怎么才能用x.startId这种语法?
如果想直接用字段名访问,需要把弱类型的Dataset[Row]转换成强类型Dataset:
- 先定义对应Schema的case class:
case class Address(endId: String, parentId: String, address: String, countries: String, sourceId: String) case class Edge(endId: String, startId: String, edgeType: String, link: String) case class Joined(endId: String, parentId: String, address: String, countries: String, sourceId: String, startId: String, edgeType: String, link: String)
- 把DataFrame转换成强类型Dataset:
val addressDS = spark.read...as[Address] val edgeDS = spark.read...as[Edge]
- join后得到强类型的
Dataset[Joined]:
val joinedDS = addressDS.join(edgeDS, "endId").as[Joined]
这时再调用groupByKey就可以直接用x.startId了:
val groupedDS = joinedDS.groupByKey(x => x.startId)
因为此时x的类型是Joined(case class),编译期就能识别startId字段,不需要依赖Spark的运行时Schema。
内容的提问来源于stack exchange,提问作者Espresso Engineer
相关产品推荐
相关产品推荐

