PySpark 3.2.x读取Parquet报INT64(TIMESTAMP(NANOS,true))错误求解
问题描述
使用PySpark提供的sqlContext.read.parquet方法读取包含时间戳列的Parquet文件时,由于PySpark写入纳秒级时间戳存在兼容性问题,选择使用pandas完成数据写入操作,pandas输出的纳秒级时间戳列格式示例:2022-05-23 08:08:35.106226000。
- PySpark 3.1.x版本可正常读取该类文件,Spark会自动将时间戳转换为long类型
- 升级到3.2.x版本后读取时抛出报错:
Illegal Parquet type: INT64 (TIMESTAMP(NANOS,true)) - 手动定义全量数据Schema进行读取的工作量过大,需要无需手动编写全量Schema的可行解决方案。
可行解决方案
按实现成本从低到高排序:
调整Spark兼容配置(优先推荐,零业务代码侵入)
Spark 3.2.x版本默认开启了纳秒级无时区时间戳(TIMESTAMP_NTZ)的严格类型校验,直接关闭该配置即可回退到3.1.x版本的解析逻辑,不需要手动定义任何字段Schema。
代码中初始化SparkSession时添加配置即可:spark = SparkSession.builder \ .config("spark.sql.parquet.timestampNTZ.enabled", "false") \ .getOrCreate()如果使用
spark-submit提交任务,直接在提交参数中追加配置:--conf spark.sql.parquet.timestampNTZ.enabled=false读取时仅转换异常时间戳字段
若不方便修改全局Spark配置,可以先读取文件拿到自动推断的Schema,只筛选出抛出异常的纳秒时间戳字段做类型转换,其余字段保持原有推断类型,不需要手动枚举全量字段:# 读取Parquet路径获取原始表结构 raw_df = spark.read.parquet("your_parquet_file_path") # 筛选出所有纳秒时间戳类型的字段 error_ts_cols = [ field.name for field in raw_df.schema if "TIMESTAMP(NANOS" in str(field.dataType) ] # 仅对异常时间戳字段做类型转换,其余字段保留 select_logic = [ f"CAST(`{col}` AS LONG) AS `{col}`" if col in error_ts_cols else f"`{col}`" for col in raw_df.columns ] result_df = raw_df.selectExpr(*select_logic)调整pandas写入侧的时间戳精度
如果可以修改写入逻辑,从根源规避兼容性问题:pandas写Parquet时将时间戳精度降到微秒级,Spark所有版本都可正常兼容读取,不需要读侧做特殊处理:# 替换为你的实际时间戳列名 df["your_timestamp_column"] = df["your_timestamp_column"].astype("datetime64[us]") df.to_parquet("your_output_parquet_path")
注意:不推荐通过降级Spark版本的方式解决该问题,上述三种方案都可以在不降级、不编写全量手动Schema的前提下解决报错。
内容的提问来源于stack exchange,提问作者Florian F.
相关产品推荐
相关产品推荐

