【发布时间】:2020-09-03 00:47:48
【问题描述】:
我在 pyspark 数据框中有数据(这是一个非常大的表,有 900M 行)
这是我拥有的数据
+-------+---------+----------+
| key| time| cond|
+-------+---------+----------+
| 6| 3704| null|
| 6| 74967| 1062|
| 6|151565068| null|
| 6|154999554| null|
| 6|160595800| null|
| 6|166192324| null|
| 6|166549533| null|
| 6|171318946| null|
| 6|754759092| null|
| 6|754999359| 18882624|
| 6|755171746| 11381128|
| 6|761097038| null|
| 6|774496554| null|
| 6|930609982| null|
| 6|930809622| null|
| 1| 192427| null|
| 1| 192427| 2779|
| 1| 717931| null|
| 1| 1110573| null|
| 1| 1155854| null|
| 1| 70049289| null|
| 1| 70687548| null|
| 1| 71222733| null|
| 1| 85006084| null|
| 1| 85029676| null|
| 1| 85032605| 1424537|
| 1| 85240114| null|
| 1| 85573757| null|
| 1| 85710915| null|
| 1| 85870370| null|
+-------+---------+----------+
这是我需要对数据框执行的操作(中间步骤):
+-------+---------+----------+--------+
| key| time| cond| result|
+-------+---------+----------+--------+
| 6| 3704| null| 0|
| 6| 74967| 1062| 1|
| 6|151565068| null| 0|
| 6|154999554| null| 1|
| 6|160595800| null| 2|
| 6|166192324| null| 3|
| 6|166549533| null| 4|
| 6|171318946| null| 5|
| 6|754759092| null| 6|
| 6|754999359| 18882624| 7|
| 6|755171746| 11381128| 0|
| 6|761097038| null| 0|
| 6|774496554| null| 1|
| 6|930609982| null| 2|
| 6|930809622| null| 3|
| 1| 192427| null| 0|
| 1| 192427| 2779| 1|
| 1| 717931| null| 0|
| 1| 1110573| null| 1|
| 1| 1155854| null| 2|
| 1| 70049289| null| 3|
| 1| 70687548| null| 4|
| 1| 71222733| null| 5|
| 1| 85006084| null| 6|
| 1| 85029676| null| 7|
| 1| 85032605| 1424537| 8|
| 1| 85240114| null| 0|
| 1| 85573757| null| 1|
| 1| 85710915| null| 2|
| 1| 85870370| null| 3|
+-------+---------+----------+--------+
'result' 列的逻辑如下:每个键都有一个运行计数器,如果 'cond' 列不为空,则将计数器归零。
我们可以假设table是orderBy("key",asc("time"))
我的最终结果实际上是条件不为空的行上的结果(每个键)的平均值。 上面的数据应该是这样的(最终结果):
+--------+--------------+
| key | avg_per_key |
+--------+--------------+
| 6| 2.66666665| ==> (1+7+0)/3
| 1| 4.5| ==> (1+8)/2
+--------+--------------+
我打算这样做:
df_results = df3[df3.cond.isNotNull()].groupby(['key']).agg(
F.expr("avg(result)").alias("avg_per_key")
)
我认为它应该可以工作,但也许有更好的方法可以在没有中间步骤的情况下做到这一点。
如何在 pyspark 中有效地做到这一点? (记住数据集很大)
【问题讨论】:
标签: pyspark