如何用pyspark.sql.functions.current_date生成指定数量的未来周日
解决PySpark中用current_date()计算未来N个周日的问题
错误原因
你写的代码报错TypeError: 'Column' object is not callable,核心问题是混淆了PySpark Column对象和Python本地datetime对象:
f.current_date()返回的是PySpark的Column类型(分布式计算的列表达式),不是Python的datetime.date实例- 你试图调用
.weekday()方法、直接和datetime.timedelta做运算,这些都是Python datetime的操作,不能直接用在Column对象上
正确解决方案(纯PySpark实现,支持打桩测试)
必须用PySpark内置的日期函数完成所有计算,避免将Column对象拉到Python本地处理,这样才能正常用@patch打桩current_date()。以下是两种实现方式:
方式1:DataFrame API实现(推荐)
from pyspark.sql import functions as f # 定义参数 N = 3 # 计算第一个即将到来的周日(若当前是周日,直接返回当天) first_sunday = f.when( f.date_format(f.current_date(), "E") == "Sun", f.current_date() ).otherwise( f.next_day(f.current_date(), "Sunday") ) # 生成N个每周间隔的周日 result_df = ( spark.range(N) # 生成0到N-1的序列 .withColumn("days_offset", f.col("id") * 7) .withColumn("Date", f.date_add(first_sunday, f.col("days_offset"))) .select("Date") ) result_df.show()
方式2:SQL表达式实现
from pyspark.sql import functions as f N = 3 # 计算第一个周日(若要包含当前周日,替换为下方注释的表达式) first_sunday = f.next_day(f.current_date(), "Sunday") # first_sunday = f.expr("CASE WHEN date_format(current_date(), 'E') = 'Sun' THEN current_date() ELSE next_day(current_date(), 'Sunday') END") # 生成日期序列并展开 result_df = ( spark.createDataFrame([(N,)], ["n"]) .withColumn("dates", f.expr(f"sequence({first_sunday}, date_add({first_sunday}, (n-1)*7), interval 7 days)")) .select(f.explode(f.col("dates")).alias("Date")) ) result_df.show()
测试验证
当current_date()被打桩为2023-05-02时,运行代码会输出:
+----------+ | Date| +----------+ |2023-05-07| |2023-05-14| |2023-05-21| +----------+
关键说明
- 所有日期计算均使用PySpark内置函数,完全兼容
current_date()的打桩需求 next_day()函数直接返回指定日期后的第一个目标星期几sequence()函数用于生成连续的日期序列,步长设为7天即可得到每周的周日- 加入了当前日期是周日的判断逻辑,确保边界情况正确处理
内容的提问来源于stack exchange,提问作者Peter
相关产品推荐
相关产品推荐

