PySpark DataFrame多行合并并聚合event_timestamp最大值的实现方法
问题描述
我有一个包含2行数据的PySpark DataFrame,需要将其合并为单行。目前已通过以下代码实现合并:
cols_to_merge = [f.first(x, ignorenulls=True).alias(x) for x in _final_df1.columns[8:23]] _final_df1.groupBy(_final_df1.columns[:8]+ _final_df1.columns[23:29]).agg(*cols_to_merge)
但现在需要在合并后的结果行中加入event_timestamp列的最大值,请问该如何实现?
原始DataFrame数据如下:
+--------------------+--------------------+----+--------------------+--------------------+---------------+-------+-----------+-------------+-------------+-------------+-------------+-------------+----------------------+----------------------+----------------------+----------------------+----------------------+----------------+----------------+----------------+----------------+----------------+------------------+------------------+------------+--------------------+-------------+--------+---------------+ | vid| sid|pvid| cid| orderId| cartId|storeId|anchor_item|substitutes_1|substitutes_2|substitutes_3|substitutes_4|substitutes_5|impressions_sub_item_1|impressions_sub_item_2|impressions_sub_item_3|impressions_sub_item_4|impressions_sub_item_5|click_sub_item_1|click_sub_item_2|click_sub_item_3|click_sub_item_4|click_sub_item_5|preferred_sub_item|preferred_sub_rank|search_query|preferred_sub_source| prefType| date_id|event_timestamp| +--------------------+--------------------+----+--------------------+--------------------+---------------+-------+-----------+-------------+-------------+-------------+-------------+-------------+----------------------+----------------------+----------------------+----------------------+----------------------+----------------+----------------+----------------+----------------+----------------+------------------+------------------+------------+--------------------+-------------+--------+---------------+ |526CC6AF-6A62-405 |911D8EC7-5B1F-4AD | |170836c4b25e4b0f9 |67374d42-024c-430 |200010858176654| 742| 14940752| null| null| null| null| null| null| null| null| null| null| 842643501| null| null| null| null| 10849171| 2| | P00N|SELECTED_PREF|20230427| 1682593258203| |526CC6AF-6A62-405 |911D8EC7-5B1F-4AD | |170836c4b25e4b0f9 |67374d42-024c-430 |200010858176654| 742| 14940752| null| null| null| null| null| null| null| null| null| null| null| 10849171| null| null| null| 10849171| 2| | P00N|SELECTED_PREF|20230427| 1682593276038| +--------------------+--------------------+----+--------------------+--------------------+---------------+-------+-----------+-------------+-------------+-------------+-------------+-------------+----------------------+----------------------+----------------------+----------------------+----------------------+----------------+----------------+----------------+----------------+----------------+------------------+------------------+------------+--------------------+-------------+--------+---------------+
当前合并后的结果如下:
+--------------------+--------------------+----+--------------------+--------------------+---------------+-------+-----------+------------------+------------------+------------+--------------------+-------------+--------+-------------+-------------+-------------+-------------+-------------+----------------------+----------------------+----------------------+----------------------+----------------------+----------------+----------------+----------------+----------------+----------------+ | vid| sid|pvid| cid| orderId| cartId|storeId|anchor_item|preferred_sub_item|preferred_sub_rank|search_query|preferred_sub_source| prefType| date_id|substitutes_1|substitutes_2|substitutes_3|substitutes_4|substitutes_5|impressions_sub_item_1|impressions_sub_item_2|impressions_sub_item_3|impressions_sub_item_4|impressions_sub_item_5|click_sub_item_1|click_sub_item_2|click_sub_item_3|click_sub_item_4|click_sub_item_5| +--------------------+--------------------+----+--------------------+--------------------+---------------+-------+-----------+------------------+------------------+------------+--------------------+-------------+--------+-------------+-------------+-------------+-------------+-------------+----------------------+----------------------+----------------------+----------------------+----------------------+----------------+----------------+----------------+----------------+----------------+ |526CC6AF-6A62-405 |911D8EC7-5B1F-4AD | |170836c4b25e4b0f9 |67374d42-024c-430 |200010858176654| 742| 14940752| 10849171| 2| | P00N|SELECTED_PREF|20230427| null| null| null| null| null| null| null| null| null| null| 842643501| 10849171| null| null| null| +--------------------+--------------------+----+--------------------+--------------------+---------------+-------+-----------+------------------+------------------+------------+--------------------+-------------+--------+-------------+-------------+-------------+-------------+-------------+----------------------+----------------------+----------------------+----------------------+----------------------+----------------+----------------+----------------+----------------+----------------+
解决方案
在原有的聚合列列表中,追加event_timestamp的最大值聚合即可,修改后的代码如下:
import pyspark.sql.functions as f cols_to_merge = [f.first(x, ignorenulls=True).alias(x) for x in _final_df1.columns[8:23]] # 添加event_timestamp的最大值聚合 cols_to_merge.append(f.max("event_timestamp").alias("event_timestamp")) _final_df1.groupBy(_final_df1.columns[:8]+ _final_df1.columns[23:29]).agg(*cols_to_merge)
说明
- 原代码通过
first(x, ignorenulls=True)合并指定列的非空值,现在新增max("event_timestamp")后,分组时会提取该分组内event_timestamp的最大值,并保留原列名。 - 执行后,合并结果中会包含
event_timestamp列,值为两行数据中的最大时间戳1682593276038。
内容的提问来源于stack exchange,提问作者Shibu
相关产品推荐
相关产品推荐

