PySpark如何在where子句中使用变量筛选DataFrame
问题根因
- 条件语句变量引用错误:你代码中的筛选条件
DPTR_DATE < 'd'里的'd'被识别为固定字符串d,而非你定义的存储日期的变量d,实际执行时是将日期类型的DPTR_DATE和字符串d做比较,没有符合条件的记录,因此返回空结果。 - Row对象未提取字段值:
d = date_range.collect()[i]取到的是PySpark的Row类型对象,而非直接的日期值,需要显式提取date字段的内容才能用于比较。
修正方案
你可以直接使用以下两种写法实现需求:
写法1:字符串拼接(简单直观)
# 一次性提取所有日期到列表,避免重复触发collect动作 date_list = [row["date"] for row in date_range.collect()] for d in date_list: # 用f-string把日期变量拼到条件语句中 df_temp = df_base.where(f"DPTR_DATE < '{d}'") display(df_temp)
写法2:API传参(更安全,避免格式错误)
from pyspark.sql import functions as F date_list = [row["date"] for row in date_range.collect()] for d in date_list: # 直接用列对象和常量值比较,不需要手动处理字符串格式 df_temp = df_base.filter(F.col("DPTR_DATE") < F.lit(d)) display(df_temp)
额外说明
如果你是在本地测试复现你给出的示例,建议把date_range的生成逻辑改成固定日期区间,避免受当前系统日期影响:
date_range = spark.sql("SELECT sequence(to_date('2021-11-13'), to_date('2021-11-15'), interval 1 day) as date").withColumn("date", F.explode(F.col("date")))
内容的提问来源于stack exchange,提问作者Daisy
相关产品推荐
相关产品推荐

