PySpark读取Parquet时日期偏移一天的问题求助
解决PySpark读取Parquet日期偏移一天的问题
核心原因分析
出现日期偏移的本质是时区不匹配或日期类型解析错误:
- pandas用
fastparquet读取时,默认将timestamp类型的日期按UTC或本地时区转换为date,而PySpark默认时区可能与pandas不一致,导致转date时偏移一天。 - 若Parquet中Date列实际存储为timestamp类型,PySpark读取后未正确处理时区,直接转date就会出错。
- 你之前的
to_date代码存在拼写错误:'yyy-MM-dd'少了一个y,这也会导致转换失败。
具体解决方案
1. 统一Spark与pandas的时区
在构建SparkSession时指定时区(如果pandas用UTC就设为UTC,用本地时区就设对应时区,比如Asia/Shanghai):
sc = SparkSession \ .builder \ .appName("test") \ .master('local[*]') \ .config("spark.sql.shuffle.partitions", "200") \ .config("spark.driver.maxResultSize","5g") \ .config("spark.executor.memory","40g") \ .config("spark.driver.memory","10g") \ .config("spark.rdd.compress", "true") \ .config("spark.sql.execution.arrow.pyspark.enabled", "true") \ .config("spark.sql.session.timeZone", "UTC") # 关键:指定时区 .getOrCreate()
2. 读取时强制指定Date列为Date类型
如果Parquet中的Date列实际是date类型,但PySpark误解析为timestamp,可通过自定义schema强制指定类型:
from pyspark.sql.types import StructType, StructField, DateType, FloatType # 根据你的表结构定义schema,替换其他列类型 schema = StructType([ StructField("Date", DateType(), nullable=True), StructField("Price", FloatType(), nullable=True) ]) df = sc.read\ .option("primitivesAsString","true")\ .option("allowNumericLeadingZeros","true")\ .schema(schema) # 应用自定义schema .parquet(f'{data_rroot}/*.parquet')
3. 正确转换Timestamp到Date(修正时区)
若Date列确实是timestamp类型,转换时需指定时区确保正确性:
import pyspark.sql.functions as F # 假设原始timestamp是UTC时区,转换为date df = df.withColumn('Date', F.to_date(F.from_utc_timestamp(df["Date"], "UTC"))) # 或者如果原始timestamp是本地时区,用to_timestamp指定时区后转date # df = df.withColumn('Date', F.to_date(F.to_timestamp(df["Date"], "yyyy-MM-dd HH:mm:ss"), "Asia/Shanghai"))
4. 调整datetimeRebaseMode参数(针对旧版pandas写入的Parquet)
如果你的Parquet文件是用pandas < 1.5.0版本写入的,需设置datetimeRebaseModeInRead为LEGACY:
sc = SparkSession \ .builder \ ... # 其他配置 .config("spark.sql.parquet.datetimeRebaseModeInRead", "LEGACY") .getOrCreate()
旧版pandas写入的Parquet日期基于Java epoch,而PySpark默认的CORRECTED模式基于公历,会导致日期偏移。
内容的提问来源于stack exchange,提问作者Peyman
相关产品推荐
相关产品推荐

