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

如何对两个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")

输出结果会是:

readerIdlocationIduserIdtime_diff_hours
R1l1u12.0
R2l1u23.0
R3l3u30.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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 06:22:03