You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 02:17:13