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

PySpark DataFrame按最大modified_time过滤失效问题求助

问题:筛选PySpark DataFrame中最新时间戳的行

我有一个包含多列的PySpark DataFrame,其中存在一个名为modified_time的时间戳列,需要筛选出modified_time为最新(即最大值)的行。

示例DataFrame

Col1  Col2  Modified_time
1     ABC   2022-03-01 08:55:29.423
2     DEF   2022-03-01 08:55:29.423
3     GHI   2022-03-01 08:55:29.423
1     ABC   2022-02-28 12:43:17.142
2     DEF   2022-02-27 23:31:26.777
3     GHI   2022-02-18 06:17:11.534 

期望输出

Col1  Col2  Modified_time
  1     ABC   2022-03-01 08:55:29.423
  2     DEF   2022-03-01 08:55:29.423
  3     GHI   2022-03-01 08:55:29.423

需求等价于:

filtered_df where modified_time = max(modified_time)

我尝试了以下代码:

df = spark.read.format("delta").table("tableName")
max_ts = df.agg({"_modified_time": "max"}).collect()[0][0]
df.filter(df._modified_time==max_ts).show(truncate = False)

但将这段代码放入整体流程中时,DataFrame并未被过滤,而是返回了所有记录,恳请帮忙解决该问题。


问题排查与解决方案

1. 优先排查列名不匹配问题

你示例中的列名是Modified_time(首字母大写),但代码中使用的是_modified_time(下划线开头),这大概率是核心问题:如果实际列名和代码引用的列名不一致,筛选条件会因无法匹配有效数据而返回全量行。先执行df.printSchema()或print(df.columns)确认DataFrame的真实列名。

2. 推荐的可靠实现方式

即使列名正确,原代码通过collect()将最大值拉到Driver端的方式,在大数据量场景下效率低且可能出现数据一致性问题,推荐以下两种更优方案:

方法一:子查询关联(避免拉取数据到Driver)

from pyspark.sql import functions as F

df = spark.read.format("delta").table("tableName")
# 生成包含最大时间戳的临时DataFrame
max_ts_df = df.select(F.max("modified_time").alias("max_ts"))
# 关联筛选出符合条件的行
filtered_df = df.join(max_ts_df, df.modified_time == max_ts_df.max_ts, "inner").drop("max_ts")
filtered_df.show(truncate=False)

方法二:窗口函数(支持全局/分组取最新)

如果后续需要扩展为分组取最新数据,窗口函数的通用性更强,全局取最新也适用:

from pyspark.sql import Window
from pyspark.sql import functions as F

df = spark.read.format("delta").table("tableName")
# 按时间戳降序排序的全局窗口
window_spec = Window.orderBy(F.desc("modified_time"))
# 生成排名列,筛选排名第一的行
filtered_df = df.withColumn("rank", F.rank().over(window_spec)) \
                .filter(F.col("rank") == 1) \
                .drop("rank")
filtered_df.show(truncate=False)

方法三:修复原代码逻辑

如果坚持使用原思路,先确保列名正确,同时优化agg的写法:

from pyspark.sql import functions as F

df = spark.read.format("delta").table("tableName")
max_ts = df.agg(F.max("modified_time")).collect()[0][0]
filtered_df = df.filter(df.modified_time == max_ts)
filtered_df.show(truncate=False)

3. 其他潜在问题排查

  • 时间戳类型问题:如果modified_time是字符串类型而非Timestamp,直接比较可能出错,先转换类型:
    df = df.withColumn("modified_time", F.to_timestamp("modified_time"))
    
  • Delta表版本问题:若Delta表有版本迭代,确保读取的是目标版本:
    df = spark.read.format("delta").option("versionAsOf", 1).table("tableName")
    

内容的提问来源于stack exchange,提问作者RLH

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 13:35:06