PySpark动态使用date_sub函数过滤DataFrame出现列解析错误求助
问题根因
你直接将Python字符串类型的application_run_date传入PySpark的to_date()函数时,该函数默认会把输入值识别为DataFrame的列名,而非你要传入的字面量日期值,因此Spark会尝试查找名为2021-09-04的列,触发列不存在的报错。
解决方案
核心是让Spark识别到传入的日期是常量而非列名,有两种实现方式可选:
方案1:Spark层使用lit()包装常量
用lit()函数将Python变量包装为Spark可识别的常量字面量,再进行日期计算,同时建议将表中的cust_date字段显式转换为日期类型,避免字符串格式不统一导致的比较错误。
from pyspark.sql.functions import * from pyspark.sql.functions import date_sub from pyspark.sql import DataFrame application_run_date = '2021-09-04' # 用lit包装字符串为Spark字面量,再转日期格式 run_date = to_date(lit(application_run_date),'yyyy-MM-dd') days_to_subtract = 1 # 先将cust_date转成日期类型再做区间过滤 final_data = raw_dataframe.filter( to_date(raw_dataframe["cust_date"], 'yyyy-MM-dd').between( date_sub(run_date, days_to_subtract), run_date ) )
方案2:Python层预处理日期
直接在Python层完成日期格式转换和区间计算,再把计算好的日期对象传入过滤逻辑:
from pyspark.sql.functions import * from datetime import datetime, timedelta from pyspark.sql import DataFrame application_run_date = '2021-09-04' days_to_subtract = 1 # Python层先算好过滤区间 run_date_dt = datetime.strptime(application_run_date, '%Y-%m-%d').date() start_date_dt = run_date_dt - timedelta(days=days_to_subtract) final_data = raw_dataframe.filter( to_date(raw_dataframe["cust_date"], 'yyyy-MM-dd').between(start_date_dt, run_date_dt) )
运行上述代码后执行final_data.show()即可得到符合过滤规则的结果:
+----+----------+--------+-----+ |cust| cust_date|purchase|value| +----+----------+--------+-----+ | 999|2021-09-03| Buy_C| 20| +----+----------+--------+-----+
内容的提问来源于stack exchange,提问作者user16798185
相关产品推荐
相关产品推荐

