【问题标题】:Transform rows to multiple rows in Spark Scala在 Spark Scala 中将行转换为多行
【发布时间】: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


    【解决方案1】:

    结果应该只有 3 列?我的意思是 col1、col2、col3(agg 结果)。 第二个聚合是不同的用户数? (我认为是的)。

    如果是这样,您基本上可以创建 3 个数据框,然后将它们合并。 阻碍:

    val df1 = spark.sql("select colA as col1, colB as col2 ,sum(Size) as colAgg group by colA,colB")

    val df2 = spark.sql("select colB as col1, colC as col2 ,Distinct(User) as colAgg group by colB,colC")

    val df3 = spark.sql("select colA as col1, colC as col2 ,sum(Size) as colAgg group by colA,colC")

    df1.union(df2).union(df3)

    【讨论】:

    • 我不想做多个联合,因为在产品中,我会有大约 200 个聚合。如果我使用数据框,也许我应该能够使用一些平面图,但我不确定如何
    猜你喜欢
    • 2015-12-05
    • 2019-02-20
    • 2017-05-20
    • 1970-01-01
    • 1970-01-01
    • 2020-08-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多