我正在嘗試使用 Pyspark 創建一個用戶保留表,我可以將其傳輸到 AWS Glue 以創建一個 ETL 作業,我可以在 QuickSight 中使用 Athena 查詢該作業。
基本上,我有兩張表,一張是用戶注冊日期,一張是用戶活動日期。然后將該注冊日期與活動日期進行比較,以計算用戶在注冊后多長時間處于活動狀態。此后,我想跟蹤在某個月份注冊的用戶中有多少在 0、1、2 周后處于活動狀態。因此,我想計算第 0 周后、第 1 周后等用戶的不同數量,即不是正常的同類群組表,其中它們按月分組然后進行跟蹤,這可能導致用戶活動在注冊后 3 個月然后在 2 個月后更大的情況。
下表和預期結果如下所示:
- user_id 1 有 5 個活動,第 0 周有 2 個,第 2 周有 2 個,第 6 周有 1 個。
- user_id 2 有 5 個活動,第 0 周有 1 個,第 1 周有 2 個,第 2 周有 1 個,第 3 周有 1 個。
- user_id 3 有 3 個活動,第 0 周有 1 個,第 1 周有 1 個,第 4 周有 1 個
然而,
- 在 8 月注冊后的第 0 周或更晚看到了 3 個唯一用戶(id:1、2、3)。
- 在 8 月注冊后 1 周或更晚看到了 3 個唯一用戶(user_id:1、2、3)。
- ...
- 在 8 月注冊后的 4 周或更晚看到了 2 個唯一用戶 (user_id: 1, 3)。
- 在 8 月注冊后的 5 周或更晚看到了 1 個唯一用戶 (user_id: 1)。
- 在 8 月注冊后的 6 周或更晚看到了 1 個唯一用戶 (user_id: 1)。
- 在 8 月注冊后 7 周或更晚看到的唯一用戶為 0。


要獲得每月的注冊數量,我只需做一個簡單的 groupBy:
df_reg = df\
.sort(col('user_id').asc(), col('created_at').asc())\
.groupBy('registered_at_month')\
.agg(countDistinct('user_id').alias('reg'))
為了在每周之后獲得不同的用戶數量,我對資料框應用了一個過濾器并回圈遍歷幾周,然后應用一個資料透視函式來獲取表格:
retention = []
for week in weeks:
print(week)
df_out = df\
.filter((col('diff_week') >= week))\
.sort(col('user_id').asc(), col('created_at').asc())\
.groupBy('registered_at_month')\
.agg(countDistinct('user_id').alias('countDistinct'))\
.withColumn('week', lit(week))
retention.append(df_out)
df_retention = functools.reduce(DataFrame.union, retention)
df_retention_2 = df_retention\
.groupBy('registered_at_month')\
.pivot('week')\
.agg(first('countDistinct'))\
.orderBy('registered_at_month')
有更清潔的方法嗎?最好沒有 for 回圈。此外,當輸入資料變大并且每月有數千名用戶注冊和數百周的活動時,樞軸函式需要永遠嗎?最后,這可以使用一些計算欄位直接在 QuickSight 中完成嗎?
非常感謝您的幫助!謝謝!
uj5u.com熱心網友回復:
是的,有一種更高效的方法可以做到這一點。在 Spark 中,按聚合分組是昂貴的,因為它意味著一個 shuffle 階段,此時 Spark 在其執行程式之間重新組織資料。在您當前的代碼中,您每周都在進行聚合,這意味著您正在執行n 2聚合,其中n是周數:一個用于注冊用戶n數,一個用于每周聚合,一個用于樞軸聚合。
您可以將其減少到兩個聚合,方法是在同一聚合中對每周求和,而不是每周求和然后進行透視。這是代碼:
from pyspark.sql import functions as F
result = df.groupby(
F.date_format('registered_at', 'MMM').alias('Month'),
F.col('user_id')
) \
.agg(F.max('diff_week').alias('max_diff')) \
.groupBy('Month') \
.agg(
F.countDistinct('user_id').alias('Registered'),
*[F.sum((F.col('max_diff') >= week).cast('integer')).alias(str(week)) for week in weeks]
) \
.orderBy('Month')
使用weeks包含從 0 到 10 的整數的陣列,以及以下df資料幀:
------------- ---------- --------- -------
|registered_at|created_at|diff_week|user_id|
------------- ---------- --------- -------
|2021-08-01 |2021-08-01|0 |1 |
|2021-08-01 |2021-08-05|0 |1 |
|2021-08-01 |2021-08-18|2 |1 |
|2021-08-01 |2021-08-21|2 |1 |
|2021-08-01 |2021-09-15|6 |1 |
|2021-08-01 |2021-08-01|0 |2 |
|2021-08-01 |2021-08-09|1 |2 |
|2021-08-01 |2021-08-10|1 |2 |
|2021-08-01 |2021-08-19|2 |2 |
|2021-08-01 |2021-08-22|3 |2 |
|2021-08-02 |2021-08-02|0 |3 |
|2021-08-02 |2021-08-09|1 |3 |
|2021-08-02 |2021-08-30|4 |3 |
------------- ---------- --------- -------
您會得到以下result輸出:
----- ---------- --- --- --- --- --- --- --- --- --- ---
|Month|Registered|0 |1 |2 |3 |4 |5 |6 |7 |8 |9 |
----- ---------- --- --- --- --- --- --- --- --- --- ---
|Aug |3 |3 |3 |3 |3 |2 |1 |1 |0 |0 |0 |
----- ---------- --- --- --- --- --- --- --- --- --- ---
它會比您的解決方案更高效
注意:在聚合之前對資料框進行排序是沒有用的,因為聚合會重新排序資料。但是,這并沒有什么壞處,因為 Spark Catalyst 在聚合之前會忽略這些排序。
轉載請註明出處,本文鏈接:https://www.uj5u.com/yidong/468336.html
