如何对两个DataFrame的列分组并对行应用聚合差值函数?
嘿,我来帮你搞定这个分组聚合差值的需求!先理清楚核心思路:我们需要先把两个DataFrame合并到一起,让同一组的记录归拢到一块,然后按照readerId、locationId、userId这几个列分组,最后对timestamp列应用差值计算逻辑。下面分Pandas和PySpark两种常用场景给你具体实现方案:
方案1:使用Pandas实现
步骤1:模拟/加载你的DataFrame数据
先把你提供的两个DataFrame转换成Pandas可处理的格式:
import pandas as pd # 第一个DataFrame df1 = pd.DataFrame({ "readerId": ["R2", "R1", "R3"], "locationId": ["l1", "l1", "l3"], "userId": ["u2", "u1", "u3"], "timestamp": pd.to_datetime(["2018-04-12 05:00:00", "2018-04-12 05:00:00", "2018-04-12 05:00:00"]) }) # 第二个DataFrame(补全你未写完的内容) df2 = pd.DataFrame({ "readerId": ["R1", "R2"], "locationId": ["l1", "l1"], "userId": ["u1", "u2"], "timestamp": pd.to_datetime(["2018-04-12 07:00:00", "2018-04-12 08:00:00"]) })
步骤2:合并两个DataFrame
把两组数据拼接到一起,确保同组记录能被后续分组捕获:
combined_df = pd.concat([df1, df2], ignore_index=True)
步骤3:分组计算时间差值
这里分两种常见需求:
- 需求A:计算每组内最早和最晚时间的差值(适合每组只有两条记录的场景)
# 分组后计算时间差(转换为小时单位,也可以用seconds/minutes等) result = combined_df.groupby(["readerId", "locationId", "userId"])["timestamp"].agg( lambda x: (x.max() - x.min()).total_seconds() / 3600 ).reset_index(name="time_diff_hours")
输出结果会是:
| readerId | locationId | userId | time_diff_hours |
|---|---|---|---|
| R1 | l1 | u1 | 2.0 |
| R2 | l1 | u2 | 3.0 |
| R3 | l3 | u3 | 0.0 |
- 需求B:计算每组内相邻时间的差值(适合每组有多条记录的场景)
先对每组内的时间戳排序,再计算相邻行的差值:
# 先按分组键+时间戳排序 combined_df_sorted = combined_df.sort_values(["readerId", "locationId", "userId", "timestamp"]) # 分组计算相邻时间差 result_diff = combined_df_sorted.groupby(["readerId", "locationId", "userId"])["timestamp"].diff().reset_index()
方案2:使用PySpark实现
如果是用Spark处理大数据量,思路类似,具体实现如下:
步骤1:模拟/加载你的DataFrame数据
from pyspark.sql import SparkSession from pyspark.sql.functions import max, min, unix_timestamp, col, lag from pyspark.sql.window import Window spark = SparkSession.builder.appName("time_diff_calculation").getOrCreate() # 第一个DataFrame df1 = spark.createDataFrame([ ("R2", "l1", "u2", "2018-04-12 05:00:00"), ("R1", "l1", "u1", "2018-04-12 05:00:00"), ("R3", "l3", "u3", "2018-04-12 05:00:00") ], ["readerId", "locationId", "userId", "timestamp"]) # 第二个DataFrame(补全内容) df2 = spark.createDataFrame([ ("R1", "l1", "u1", "2018-04-12 07:00:00"), ("R2", "l1", "u2", "2018-04-12 08:00:00") ], ["readerId", "locationId", "userId", "timestamp"])
步骤2:合并并转换时间格式
# 合并两个DataFrame combined_df = df1.union(df2) # 将字符串类型的timestamp转换为Spark时间类型 combined_df = combined_df.withColumn("timestamp", col("timestamp").cast("timestamp"))
步骤3:分组计算时间差值
- 需求A:每组最早最晚时间差
result = combined_df.groupBy("readerId", "locationId", "userId") \ .agg( (unix_timestamp(max("timestamp")) - unix_timestamp(min("timestamp")))/3600 \ .alias("time_diff_hours") ) result.show()
输出结果:
+--------+----------+------+---------------+ |readerId|locationId|userId|time_diff_hours| +--------+----------+------+---------------+ | R2| l1| u2| 3.0| | R1| l1| u1| 2.0| | R3| l3| u3| 0.0| +--------+----------+------+---------------+
- 需求B:每组相邻时间差
用窗口函数lag获取上一条记录的时间,再计算差值:
# 定义窗口:按分组键分区,按时间戳排序 window_spec = Window.partitionBy("readerId", "locationId", "userId").orderBy("timestamp") # 计算相邻时间差(小时单位) result_diff = combined_df.withColumn( "prev_timestamp", lag("timestamp", 1).over(window_spec) ).withColumn( "time_diff_hours", (unix_timestamp("timestamp") - unix_timestamp("prev_timestamp"))/3600 ) result_diff.show()
内容的提问来源于stack exchange,提问作者Vinju
相关产品推荐
相关产品推荐

