如果不同组中的数组大小不等(对您而言,组是("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] |
+-----+----------------------------------------+