Spark读取Parquet文件遇未知13位日期数据类型的解决求助
解决Parquet文件中Hive/Impala生成的int96日期列识别问题
一、先确认列的原始存储类型
先通过Parquet元数据工具确认目标日期列确实是Hive/Impala常用的int96类型,执行命令:
parquet-tools meta your_file.parquet
查看输出中目标列的type字段,确认是否为INT96。
二、Spark环境下的解决方案
如果使用Spark读取,除了设置重基参数,还需注意版本兼容性和显式Schema的正确绑定:
- 基础配置(Spark 3.0+)
确保启用LEGACY重基模式的同时,显式指定列类型为TimestampType:val spark = SparkSession.builder() .config("spark.sql.parquet.int96RebaseModeInRead", "LEGACY") .config("spark.sql.parquet.datetimeRebaseModeInRead", "LEGACY") .getOrCreate() import org.apache.spark.sql.types._ val customSchema = StructType(Seq( StructField("target_date_col", TimestampType, nullable = true), // 其他列定义... )) val df = spark.read.schema(customSchema).parquet("path/to/your/parquet/file") - 手动字节解析(兜底方案)
如果上述配置无效,直接解析int96的原始字节(Hive的int96是4字节Julian日+8字节当天纳秒,大端字节序):import java.time._ import java.nio.ByteBuffer import java.nio.ByteOrder val parseInt96Timestamp = udf((bytes: Array[Byte]) => { if (bytes == null || bytes.length != 12) null else { // 解析Julian日(前4字节) val julianDay = ByteBuffer.wrap(bytes.take(4)).order(ByteOrder.BIG_ENDIAN).getInt() // 解析当天纳秒数(后8字节) val nanosOfDay = ByteBuffer.wrap(bytes.drop(4)).order(ByteOrder.BIG_ENDIAN).getLong() // Julian日转公历:基准日为公元前4713年11月24日 val baseDate = LocalDate.of(-4713, 11, 24) val targetDate = baseDate.plusDays(julianDay) // 组合为完整时间戳 Timestamp.valueOf(targetDate.atStartOfDay().plusNanos(nanosOfDay)) } }) val parsedDF = df.withColumn("parsed_date", parseInt96Timestamp(col("target_date_col")))
三、Python/Pandas环境下的解决方案
- PyArrow引擎配置
使用PyArrow读取时,指定int96的解析规则:import pandas as pd df = pd.read_parquet( "path/to/your/file.parquet", engine="pyarrow", pyarrow_options={ "int96_timestamp_unit": "ns", "timestamp_parquet_version": "1.0" # 匹配Hive/Impala的Parquet版本 } ) - 手动字节解析
直接处理原始字节列:import pyarrow.parquet as pq import datetime table = pq.read_table("path/to/your/file.parquet") date_bytes = table.column("target_date_col").to_pylist() parsed_dates = [] julian_base = datetime.date(year=-4713, month=11, day=24) for b in date_bytes: if b is None: parsed_dates.append(None) continue # 解析Julian日和纳秒数 julian_day = int.from_bytes(b[:4], byteorder='big', signed=True) nanos = int.from_bytes(b[4:], byteorder='big', signed=True) # 计算完整日期时间 target_date = julian_base + datetime.timedelta(days=julian_day) target_datetime = datetime.datetime.combine(target_date, datetime.time.min) + datetime.timedelta(nanoseconds=nanos) parsed_dates.append(target_datetime) # 替换原列 df = table.to_pandas() df["target_date_col"] = parsed_dates
四、核心转换逻辑说明
Hive/Impala的int96时间戳存储规则是:
- 12字节分为两部分:前4字节是Julian日数(从公元前4713年11月24日开始计算),后8字节是当天的纳秒偏移量
- 转换时先把Julian日转成公历日期,再加上纳秒偏移量得到完整时间戳
内容的提问来源于stack exchange,提问作者Code Noober
相关产品推荐
相关产品推荐

