【发布时间】:2020-05-30 19:25:28
【问题描述】:
我有一个问题,我需要将一行转换为多行。这是基于我拥有的不同映射。我试图在下面提供一个示例。
假设我有一个具有以下架构的镶木地板文件
ColA, ColB, ColC, Size, User
我需要根据查找图将上述数据聚合成多行。假设我有一张静态地图
ColA, ColB, Sum(Size)
ColB, ColC, Distinct (User)
ColA, ColC, Sum(Size)
这意味着输入 RDD 中的一行需要转换为 3 个聚合。我相信 RDD 是使用 FlatMapPair 的方式,但我不确定如何去做。
我也可以将列连接成一个键,例如ColA_ColB 等。
为了从相同的数据创建多个聚合,我从这样的东西开始
val keyData: PairFunction[Row, String, Long] = new PairFunction[Row, String, Long]() {
override def call(x: Row) = {
(x.getString(1),x.getLong(5))
}
}
val ip15M = spark.read.parquet("a.parquet").toJavaRDD
val pairs = ip15M.mapToPair(keyData)
java.util.List[(String, Long)] = [(ios,22), (ios,23), (ios,10), (ios,37), (ios,26), (web,52), (web,1)]
我相信我需要做 flatmaptopair 而不是 mapToPair。在类似的线路上,我尝试过
val FlatMapData: PairFlatMapFunction[Row, String, Long] = new PairFlatMapFunction[Row, String, Long]() {
override def call(x: Row) = {
(x.getString(1),x.getLong(5))
}
}
但它给出了错误
Expression of type (String, Long) doesn't conform to expected type util.Iterator[(String, Long)]
感谢任何帮助。如果我需要添加更多详细信息,请告诉我。
【问题讨论】:
标签: scala apache-spark rdd