PySpark生成指定起止日期的日期DataFrame及报错排查
生成日期序列DataFrame的PySpark解决方案
问题说明
给定起始日期start_date = "1900-01-01",结束日期end_date = "2022-01-31",希望用PySpark生成包含该时间段内所有日期的DataFrame,但运行以下代码时触发错误:
int() argument must be a string, a bytes-like object or a number, not 'Column'
原代码如下:
from pyspark.sql.functions import col, date_add, to_date from pyspark.sql import SparkSession # 创建SparkSession spark = SparkSession.builder.appName("DateDataFrame").getOrCreate() # 定义起始和结束日期 start_date = "1900-01-01" end_date = "2022-01-31" date_df = spark.range(0, (to_date(col(end_date)) - to_date(col(start_date))).cast("int")).select(date_add(to_date(start_date), col("id")).alias("date"))
错误原因
spark.range()的上下限必须是Python原生数值,但原代码中(to_date(col(end_date)) - to_date(col(start_date))).cast("int")返回的是PySpark的Column对象,无法直接作为range()的参数。col(end_date)用法错误:end_date是Python字符串变量,不是DataFrame的列名,不需要用col()包裹。
修正后的代码
方式一:用PySpark计算日期差
先在Driver端计算出起始日到结束日的总天数,再用该数值生成序列,最后逐个生成日期:
from pyspark.sql.functions import date_add, to_date from pyspark.sql import SparkSession from pyspark.sql.types import DateType # 创建SparkSession spark = SparkSession.builder.appName("DateDataFrame").getOrCreate() # 定义起始和结束日期 start_date = "1900-01-01" end_date = "2022-01-31" # 计算起始日到结束日的总天数(+1是为了包含结束日期) days_diff = spark.sql(f"SELECT datediff('{end_date}', '{start_date}') + 1 AS diff").collect()[0][0] # 生成日期序列DataFrame date_df = spark.range(0, days_diff).select(date_add(to_date(start_date), col("id")).alias("date")) # 查看前5条结果 date_df.show(5)
方式二:用Python datetime计算日期差(更简洁)
直接用Python原生datetime模块计算天数差,再传入spark.range():
from pyspark.sql.functions import date_add, to_date from pyspark.sql import SparkSession from datetime import datetime spark = SparkSession.builder.appName("DateDataFrame").getOrCreate() start_date = "1900-01-01" end_date = "2022-01-31" # 转换为datetime对象并计算天数差 start_dt = datetime.strptime(start_date, "%Y-%m-%d") end_dt = datetime.strptime(end_date, "%Y-%m-%d") days_diff = (end_dt - start_dt).days + 1 # 生成日期序列 date_df = spark.range(0, days_diff).select(date_add(to_date(start_date), col("id")).alias("date")) date_df.show(5)
内容的提问来源于stack exchange,提问作者john
相关产品推荐
相关产品推荐

