如何避免PySpark读取Pandas datetime64[ns]生成的Parquet时的分析异常
解决Pandas生成的含datetime64[ns]列的Parquet文件在Spark中读取时的类型异常问题
问题描述
当读取由含datetime64[ns]列的Pandas DataFrame生成的Parquet文件时,Spark会抛出如下异常:
org.apache.spark.sql.AnalysisException: Illegal Parquet type: INT64 (TIMESTAMP(NANOS,false))
复现代码:
import pandas as pd from pyspark.sql import SparkSession # 生成含datetime64[ns]列的Pandas DataFrame并写入Parquet pdf = pd.DataFrame(data={'time': pd.date_range('10/1/23', '10/7/23', freq='D')}) pdf.to_parquet('<path>/data.parquet') # Spark读取Parquet时触发异常 spark = SparkSession.builder.getOrCreate() sdf = spark.read.parquet('<path>/data.parquet') # 尝试转换回Pandas DataFrame df = sdf.toPandas()
目标是让Spark读取该Parquet文件时,将time列正确识别为TimestampType,同时避免现有非理想方案(如将列转为object类型)带来的FutureWarning。
解决方案
方案1:修改Pandas写入Parquet的参数,适配Spark时间戳格式
Pandas默认以纳秒精度的INT64类型存储datetime64[ns]列到Parquet,而Spark对该格式兼容性较差。可以使用pyarrow引擎,并指定时间戳的写入格式为Spark兼容的类型:
# 使用pyarrow引擎写入,指定时间戳格式为timestamp[ns] pdf.to_parquet('<path>/data.parquet', engine='pyarrow', write_timestamp='timestamp[ns]')
之后Spark读取时即可正确识别为TimestampType。
方案2:配置Spark以支持纳秒精度时间戳
在创建SparkSession时添加配置参数,开启对纳秒精度时间戳的支持:
spark = SparkSession.builder \ .config("spark.sql.timestampType", "TIMESTAMP_NS") \ .config("spark.sql.parquet.datetimeRebaseModeInRead", "CORRECTED") \ .getOrCreate() # 读取Parquet文件 sdf = spark.read.parquet('<path>/data.parquet')
该配置让Spark能够直接解析Pandas写入的纳秒精度时间戳列。
方案3:转换Pandas时间戳为毫秒精度(业务允许时)
如果业务不需要纳秒级精度,可以将Pandas的datetime64[ns]列转换为毫秒精度后再写入Parquet,Spark可直接识别:
pdf['time'] = pdf['time'].astype('datetime64[ms]') pdf.to_parquet('<path>/data.parquet')
避坑提示
不要使用将datetime64[ns]列转为object类型的方案(如pdf['time'] = pd.Series(pdf['time'].dt.to_pydatetime(), dtype=object)),该方法会在后续将Spark DataFrame转换回Pandas时触发FutureWarning,且不符合类型规范。
内容的提问来源于stack exchange,提问作者Russell Burdt
相关产品推荐
相关产品推荐

