【问题标题】:Aggregating several fields simultaneously from Dataset从数据集中同时聚合多个字段
【发布时间】:2018-07-10 06:55:11
【问题描述】:

我有一个具有以下方案的数据:

sourceip
destinationip
packets sent

我想从这些数据中计算出几个聚合字段并具有以下架构:

ip 
packets sent as sourceip
packets sent as destination

在 RDD 的快乐日子里,我可以使用 aggregate,定义 {ip -> []} 的映射,并在相应的数组位置计算出现次数。

在数据集/数据框聚合不再可用,而是可以使用 UDAF,不幸的是,根据我使用 UDAF 的经验,它们是不可变的,这意味着它们不能被使用(必须在每次地图更新时创建一个新实例) example + explanation here

一方面,从技术上讲,我可以将数据集转换为 RDD、聚合等,然后返回数据集。我预计这会导致性能下降,因为数据集更加优化。由于复制,UDAF 是不可能的。

还有其他方法可以执行聚合吗?

【问题讨论】:

    标签: apache-spark apache-spark-sql


    【解决方案1】:

    听起来你需要一个标准的melt (How to melt Spark DataFrame?) 和pivot 组合:

    val df = Seq(
      ("192.168.1.102", "192.168.1.122", 10),
      ("192.168.1.122", "192.168.1.65", 10),
      ("192.168.1.102", "192.168.1.97", 10)
    ).toDF("sourceip", "destinationip", "packets sent")
    
    
    df.melt(Seq("packets sent"), Seq("sourceip", "destinationip"), "type", "ip")
      .groupBy("ip")
      .pivot("type", Seq("sourceip", "destinationip"))
      .sum("packets sent").na.fill(0).show
    
    // +-------------+--------+-------------+             
    // |           ip|sourceip|destinationip|
    // +-------------+--------+-------------+
    // | 192.168.1.65|       0|           10|
    // |192.168.1.102|      20|            0|
    // |192.168.1.122|      10|           10|
    // | 192.168.1.97|       0|           10|
    // +-------------+--------+-------------+
    

    【讨论】:

    • 很好地实现了 column-wise-explode + pivot 用法。 thx alot 完美地工作
    【解决方案2】:

    在没有任何自定义聚合的情况下进行此操作的一种方法是使用 flatMap(或 explode 用于数据帧),如下所示:

    case class Info(ip : String, sent : Int, received : Int)
    case class Message(from : String, to : String, p : Int)
    val ds = Seq(Message("ip1", "ip2", 5), 
                 Message("ip2", "ip3", 7), 
                 Message("ip2", "ip1", 1), 
                 Message("ip3", "ip2", 3)).toDS()
    
    ds
        .flatMap(x => Seq(Info(x.from, x.p, 0), Info(x.to, 0, x.p)))
        .groupBy("ip")
        .agg(sum('sent) as "sent", sum('received) as "received")
        .show
    
    
    // +---+----+--------+
    // | ip|sent|received|
    // +---+----+--------+
    // |ip2|   8|       8|
    // |ip3|   3|       7|
    // |ip1|   5|       1|
    // +---+----+--------+
    

    就性能而言,我不确定 flatMap 与自定义聚合相比是否有所改进。

    【讨论】:

    • 非常感谢!这似乎是一个完美的解决方案。我不得不接受 user8371915 的解决方案,不是因为你的解决方案更糟,而是因为 flatMap 在 pyspark 中不可用(抱歉,不得不提的是前面的)。
    • 没问题,很高兴您找到了答案。如果您有兴趣,我添加了另一个具有相同逻辑的答案,仅使用 pyspark 数据帧。
    【解决方案3】:

    这是一个使用 explode 的 pyspark 版本。比较冗长但是逻辑和flatMap版本完全一样,只是纯dataframe代码。

    sc\
      .parallelize([("ip1", "ip2", 5), ("ip2", "ip3", 7), ("ip2", "ip1", 1), ("ip3", "ip2", 3)])\
      .toDF(("from", "to", "p"))\
      .select(F.explode(F.array(\
          F.struct(F.col("from").alias("ip"),\
                   F.col("p").alias("received"),\
                   F.lit(0).cast("long").alias("sent")),\
          F.struct(F.col("to").alias("ip"),\
                   F.lit(0).cast("long").alias("received"),\
                   F.col("p").alias("sent")))))\
      .groupBy("col.ip")\
      .agg(F.sum(F.col("col.received")).alias("received"), F.sum(F.col("col.sent")).alias("sent"))
    
    // +---+----+--------+
    // | ip|sent|received|
    // +---+----+--------+
    // |ip2|   8|       8|
    // |ip3|   3|       7|
    // |ip1|   5|       1|
    // +---+----+--------+
    

    【讨论】:

    • 第一,谢谢时间的投入,逻辑是不是和user8371915提供的melt函数一样?
    • 这个方法非常相似,因为melt 也会爆炸数据框。然而,这种方法不需要额外的支点,并且可能会更有效。
    • 您能否详细说明为什么它会更有效率?我猜是因为您手动定义了要聚合的字段,而不是使用通用的pivot,您能同意吗?
    • Melt 生成类型为 (ip, string, int) 的记录,而这个生成类型为 (ip, int, int) 的记录将占用更少的空间。这可以通过用索引替换字符串来解决。此外,熔化/旋转/填充稍微复杂一些(它产生 3 个火花阶段,而这一阶段为 2 个,2 个随机阶段对 1 个)。在我生成的一些样本数据中,融合/枢轴方法需要多花大约 50% 的时间。
    【解决方案4】:

    由于您没有提及上下文和聚合,您可以执行以下操作,

    val df = ??? // your dataframe/ dataset
    

    来自 Spark 来源:

    (Scala-specific)通过指定列中的映射来计算聚合 聚合方法的名称。生成的 DataFrame 还将包含 分组列。可用的聚合方法有 avg、max、 分,总和,计数。

    // 选择最年长员工的年龄和总费用 每个部门

     df
     .groupBy("department")
     .agg(Map(
          "age" -> "max",
          "expense" -> "sum"   
         ))
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-11-08
      • 2021-12-16
      • 2015-11-26
      • 2018-10-03
      • 1970-01-01
      • 2022-09-23
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多