DBR 12.2 LTS中Spark 3.3.2读取含INT64 TIMESTAMP[NANOS]的Parquet失败
在挂载至DBR 12.2 LTS(内置Apache Spark 3.3.2、Scala 2.12)的Databricks作业集群Notebook中,执行代码spark.read.parquet("dbfs:/path/test.parquet")在2024-08-09前可正常运行,之后突然失效。
目标Parquet文件的Schema为:
x: DOUBLE, OPTIONAL valid_time: INT64, TIMESTAMP[NANOS], OPTIONAL
排查结论:
- 移除
valid_time字段的元数据后,Spark读取代码可正常执行,说明问题出在Spark对带有INT64类型+TIMESTAMP[NANOS]元数据的Parquet文件的处理上 - 使用Pandas代码
import pandas as pd; pd.read_parquet('/dbfs/path/test.parquet', engine='fastparquet')可正常打开文件,证明文件本身未损坏
已完成的排查动作:
- 确认Apache JIRA的SPARK-40819问题与当前场景完全匹配,但该问题标注为在Spark 3.3.2版本已修复,而当前正使用该版本
- 尝试配置Spark 3.3.2的相关参数,问题仍未解决,配置代码如下:
spark.conf.set("spark.sql.legacy.parquet.int96RebaseModeInRead", "CORRECTED") spark.conf.set("spark.sql.legacy.parquet.datetimeRebaseModeInRead", "CORRECTED") spark.conf.set("spark.sql.parquet.int96AsTimestamp", "true") spark.conf.set("spark.sql.parquet.int96TimestampConversion", "true") spark.conf.set("spark.sql.parquet.outputTimestampType", "TIMESTAMP_MICROS") spark.conf.set("spark.sql.parquet.int64AsTimestamp", "true") spark.conf.set("spark.sql.parquet.int64TimestampConversion", "true")
问题1:为何2024-08-09后突然失效?
最可能的原因是Databricks在2024-08-09前后对DBR 12.2 LTS集群的底层Spark依赖或相关组件进行了静默更新,导致原本修复SPARK-40819的补丁被移除、覆盖,或者引入了新的兼容性问题。
虽然SPARK-40819标注为在Spark 3.3.2中修复,但Databricks的DBR版本是经过定制的发行版,并非原生Apache Spark。这种定制可能包含Databricks自己的补丁、组件替换或版本微调,当Databricks推送集群维护更新时,可能意外破坏了原本对INT64类型Timestamp的处理逻辑。
另一种可能是集群的Parquet相关依赖库(如parquet-hadoop)版本被更新,新的库版本与Spark 3.3.2的处理逻辑产生冲突,导致原本支持的INT64+TIMESTAMP[NANOS]解析逻辑失效。
问题2:如何解决该问题(优先保留Spark性能)
方案1:强制指定Schema,绕过自动元数据解析
Spark自动解析Parquet元数据时会识别TIMESTAMP[NANOS]注解并尝试用INT64转Timestamp的逻辑处理,我们可以手动指定Schema,将valid_time先读取为Long类型,再转换为Timestamp,避免触发有问题的自动解析逻辑:
from pyspark.sql.types import StructType, StructField, DoubleType, LongType, TimestampType # 定义手动Schema,先将valid_time设为LongType custom_schema = StructType([ StructField("x", DoubleType(), nullable=True), StructField("valid_time", LongType(), nullable=True) ]) # 读取文件 df = spark.read.schema(custom_schema).parquet("dbfs:/path/test.parquet") # 将Long类型的纳秒时间戳转换为Timestamp df = df.withColumn("valid_time", (df["valid_time"] / 1e9).cast(TimestampType()))
这种方式完全使用Spark原生API,不会损失性能,同时绕过了自动解析时的兼容性问题。
方案2:修复Parquet文件元数据(一次性操作)
如果目标文件是固定的,可先通过Pandas读取后重新写入为Spark兼容的Parquet格式,后续直接用Spark读取新文件:
import pandas as pd # 读取原文件 pd_df = pd.read_parquet('/dbfs/path/test.parquet', engine='fastparquet') # 重新写入为Spark兼容的Parquet格式 pd_df.to_parquet('/dbfs/path/test_fixed.parquet', engine='pyarrow') # 后续用Spark读取修复后的文件 df = spark.read.parquet('/dbfs/path/test_fixed.parquet')
此方案只需执行一次,后续读取可完全保留Spark性能优势。
方案3:调整DBR版本
如果确认是DBR 12.2 LTS的静默更新导致问题,可以尝试:
- 切换到DBR 12.2 LTS的历史版本(若Databricks提供回滚选项),找到2024-08-09前的可用版本
- 升级到更高版本的DBR(如DBR 13.x或14.x LTS),这些版本可能重新修复了该兼容性问题
方案4:提交Databricks支持工单
由于问题是DBR定制版本导致的,直接联系Databricks技术支持,提供文件样本和复现步骤,要求其修复该兼容性问题,这是最彻底的长期解决方案。
内容的提问来源于stack exchange,提问作者Kelly Ma

