【问题标题】:PySpark: how to aggregate over column arrays with variable width?PySpark:如何聚合具有可变宽度的列数组?
【发布时间】:2019-11-22 18:47:53
【问题描述】:

我正在尝试聚合并创建一个方法数组(这是一个最小的工作示例):

n = len(allele_freq_total.select("alleleFrequencies").first()[0])

allele_freq_by_site = allele_freq_total.groupBy("contigName", "start", "end", "referenceAllele").agg(
  array(*[mean(col("alleleFrequencies")[i]) for i in range(n)]).alias("mean_alleleFrequencies")

使用我从

获得的解决方案

Aggregate over column arrays in DataFrame in PySpark?

但问题是n 是可变的,我该如何改变

array(*[mean(col("alleleFrequencies")[i]) for i in range(n)])

所以它考虑到可变长度?

【问题讨论】:

    标签: python-3.x apache-spark pyspark


    【解决方案1】:

    如果不同组中的数组大小不等(对您而言,组是("contigName", "start", "end", "referenceAllele"),我将简单地重命名为group),您可以考虑使用爆炸数组列(alleleFrequencies)引入值在数组中的位置。这将为您提供一个额外的列,您可以在分组中使用它来计算您想到的平均值。此时,您实际上可能有足够的空间进行进一步计算(请参阅下面的df3.show())。

    如果你真的必须把它放回数组中,那就更难了,我不知道。必须跟踪订单,我相信使用地图(如果您愿意,可以使用字典)很容易。为此,我在两列上使用聚合函数collect_list。虽然collect_list 不是确定性的(您不知道值在列表中返回的顺序,因为行被打乱了),但两个数组上的聚合将保留它们的顺序,因为行被整体打乱(见下文df4.show())。从那里,您可以使用map_from_arrays 创建位置到平均值的映射。

    例子:

    >>> from pyspark.sql.functions import mean, col, posexplode, collect_list, map_from_arrays
    >>> 
    >>> df = spark.createDataFrame([
    ...     ("A", [0, 1, 2]),
    ...     ("A", [0, 3, 6]),
    ...     ("B", [1, 2, 4, 5]),
    ...     ("B", [1, 2, 6, 1])],
    ...     schema=("group", "values"))
    >>> df2 = df.select(df.group, posexplode(df.values))  # adds the "pos" and "col" columns
    >>> df3 = (df2
    ...        .groupBy("group", "pos")
    ...        .agg(mean(col("col")).alias("avg_of_positions"))
    ...        )
    >>> df4 = (df3
    ...        .groupBy("group")
    ...        .agg(
    ...          collect_list("pos").alias("pos"),
    ...          collect_list("avg_of_positions").alias("avgs")
    ...          )
    ...        )
    >>> df5 = df4.select(
    ...     "group",
    ...     map_from_arrays(col("pos"), col("avgs")).alias("positional_averages")
    ... )
    >>> df5.show(truncate=False)
    [Stage 0:>                                                          (0 + 4) / 4]
    +-----+----------------------------------------+                                
    |group|positional_averages                     |
    +-----+----------------------------------------+
    |B    |[0 -> 1.0, 1 -> 2.0, 3 -> 3.0, 2 -> 5.0]|
    |A    |[0 -> 0.0, 1 -> 2.0, 2 -> 4.0]          |
    +-----+----------------------------------------+
    

    【讨论】:

      猜你喜欢
      • 2018-06-01
      • 2021-01-14
      • 1970-01-01
      • 2016-12-23
      • 1970-01-01
      • 2023-03-13
      • 1970-01-01
      • 1970-01-01
      • 2022-09-28
      相关资源
      最近更新 更多