【问题标题】:Spark aggregations where output columns are functions and rows are columns输出列是函数,行是列的 Spark 聚合
【发布时间】:2020-04-07 20:27:28
【问题描述】:

我想在数据帧的不同列上计算一堆不同的聚合函数。

我知道我可以做这样的事情,但输出都是一行。

df.agg(max("cola"), min("cola"), max("colb"), min("colb"))

假设我将在 10 个不同的列上执行 100 个不同的聚合。

我希望输出数据框是这样的 -

      |Min|Max|AnotherAggFunction1|AnotherAggFunction2|...etc..
cola  | 1 | 10| ... 
colb  | 2 | NULL| ... 
colc  | 5 | 20| ... 
cold  | NULL | 42| ... 
...

我的行是我正在对其执行聚合的每一列,我的列是聚合函数。例如,如果我不计算 colb max,某些区域将为空。

我怎样才能做到这一点?

【问题讨论】:

    标签: python apache-spark pyspark pyspark-sql pyspark-dataframes


    【解决方案1】:

    您可以创建一个 Map 列,例如 Metrics,其中键是列名,值是聚合结构(max、min、avg...)。我正在使用map_from_entries 函数创建地图列(可从 Spark 2.4+ 获得)。然后,只需分解地图即可获得所需的结构。

    这是一个您可以根据自己的要求进行调整的示例:

    df = spark.createDataFrame([("A", 1, 2), ("B", 2, 4), ("C", 5, 6), ("D", 6, 8)], ['cola', 'colb', 'colc'])
    
    agg = map_from_entries(array(
        *[
            struct(lit(c),
                   struct(max(c).alias("Max"), min(c).alias("Min"))
                   )
            for c in df.columns
        ])).alias("Metrics")
    
    df.agg(agg).select(explode("Metrics").alias("col", "Metrics")) \
        .select("col", "Metrics.*") \
        .show()
    
    #+----+---+---+
    #|col |Max|Min|
    #+----+---+---+
    #|cola|D  |A  |
    #|colb|6  |1  |
    #|colc|8  |2  |
    #+----+---+---+
    

    【讨论】:

      【解决方案2】:

      这是一种允许您从预定义列表动态设置聚合的解决方案。该解决方案使用 map_from_arrays 等,因此与 Spark >= 2.4.0 兼容:

      from pyspark.sql.functions import lit, expr, array, map_from_arrays
      
      df = spark.createDataFrame([
        [1, 2.3, 5000],
        [2, 5.3, 4000],
        [3, 2.1, 3000],
        [4, 1.5, 4500]
      ], ["cola", "colb", "colc"])
      
      aggs = ["min", "max", "avg", "sum"]
      aggs_select_expr = [f"value[{idx}] as {agg}" for idx, agg in enumerate(aggs)]
      
      agg_keys = []
      agg_values = []
      
      # generate map here where key is col name and value an array of aggregations
      for c in df.columns:
        agg_keys.append(lit(c)) # the key i.e cola
        agg_values.append(array(*[expr(f"{a}({c})") for a in aggs])) # the value i.e [expr("min(a)"), expr("max(a)"), expr("avg(a)"), expr("sum(a)")]
      
      df.agg(
        map_from_arrays(
          array(agg_keys), 
          array(agg_values)
        ).alias("aggs")
      ) \
      .select(explode("aggs")) \
      .selectExpr("key as col", *aggs_select_expr) \
      .show(10, False)
      
      # +----+------+------+------+-------+
      # |col |min   |max   |avg   |sum    |
      # +----+------+------+------+-------+
      # |cola|1.0   |4.0   |2.5   |10.0   |
      # |colb|1.5   |5.3   |2.8   |11.2   |
      # |colc|3000.0|5000.0|4125.0|16500.0|
      # +----+------+------+------+-------+
      

      说明: 使用表达式array(*[expr(f"{a}({c})") for a in aggs]) 我们创建一个包含当前列的所有聚合的数组。生成的数组的每个项目都使用语句expr(f"{a}({c})" 进行评估,这将产生即expr("min(a)")

      该数组将包含agg_values 的值,它们与agg_keys 将通过表达式map_from_arrays(array(agg_keys), array(agg_values)) 组成我们的最终映射。 map的结构是这样的:

      map(
          cola -> [min(cola), max(cola), avg(cola), sum(cola)]
          colb -> [min(colb), max(colb), avg(colb), sum(colb)]
          colc -> [min(cola), max(colc), avg(cola), sum(colc)]
      )
      

      为了提取我们需要的信息,我们必须用explode("aggs") 分解先前的地图,这将创建两列keyvalue,我们在选择语句中使用它们。

      aggs_select_expr 将包含["value[0] as min", "value[1] as max", "value[2] as avg", "value[3] as sum"] 形式的值,这将是selectExpr statememnt 的输入。

      更新:

      我意识到通过省略聚合(也称为隐式 groupByagg)有一种更高效的方法。我们可以通过create_map 内置函数实现同样的功能:

      from pyspark.sql.functions import create_map, expr, array
      from itertools import chain
      
      df = spark.createDataFrame([
        [1, 2.3, 5000],
        [2, 5.3, 4000],
        [3, 2.1, 3000],
        [4, 1.5, 4500]
      ], ["cola", "colb", "colc"])
      
      aggs = ["min", "max", "avg", "sum"]
      aggs_select_expr = [f"value[{idx}] as {agg}" for idx, agg in enumerate(aggs)]
      
      df.select(explode(
                        create_map(*list(
                              chain(*[(lit(c), array(*[expr(f"{a}({c})") for a in aggs])) 
                                for c in df.columns
                           ])))
                         )
               ) \
              .selectExpr("key as col", *aggs_select_expr)
      

      注意:除了代码更少之外,第二种方法的主要优点是它只包含窄转换而不包含宽转换,即groupBy。我们将提高性能,因为它避免了洗牌。

      【讨论】:

        猜你喜欢
        • 2019-06-02
        • 2016-02-26
        • 2014-09-08
        • 2019-02-22
        • 1970-01-01
        • 2021-04-07
        • 1970-01-01
        • 2011-09-05
        • 1970-01-01
        相关资源
        最近更新 更多