Spark生成Parquet时DoubleType转DecimalType(10,6)报错排查
问题描述
我有如下JSON格式数据:
{"C_0_0": "INDIA", "C_0_1": 1219.765432, "C_0_2": {"C_1_0": "INDIA", "C_1_1": ["INDIA", "INDIA"]}}
希望将其保存为Parquet文件,但默认情况下C_0_1列会被识别为DoubleType,需要转换为DecimalType(10,6)。采用以下代码实现后,生成的Parquet文件无法解析C_0_1列,报错:pyarrow.lib.ArrowNotImplementedError: Cannot extract statistics for type。
代码如下:
sc = SparkContext.getOrCreate() hc = HiveContext(sc) sqlContext = SQLContext(sc) rdd = sc.parallelize([json_str]) nested_df = hc.read.json(rdd) nested_df.show(20,False) print("Original Pyspark Schema: ", nested_df.schema) new_pyspark_schema = StructType([ StructField('C_0_0', StringType(), True), StructField('C_0_1', DecimalType(10,6), True), StructField('C_0_2', StructType([ StructField('C_1_0', StringType(), True), StructField('C_1_1', ArrayType(StringType(), True), True) ]), True ) ]) new_rdd = sc.parallelize([json_str]) new_df = hc.read.schema(new_pyspark_schema).json(new_rdd) new_df.repartition(1).write.option("schema", new_pyspark_schema).parquet(file_location) parDf = sqlContext.read.parquet(file_location) parDf.show(20,False) parDf.printSchema()
原因分析与解决方法
核心问题
- 直接指定DecimalType读取JSON的兼容性问题:Spark读取JSON时,直接用
DecimalType定义schema,JSON中的浮点数值无法被正确映射为Decimal类型,导致Parquet写入时类型逻辑混乱,后续读取触发Arrow统计信息提取失败。 - 写入Parquet时冗余的schema参数:
write.option("schema", new_pyspark_schema)是无效配置,Parquet写入会直接沿用DataFrame的schema,这个多余参数可能干扰类型处理流程。
修正方案
不要直接用DecimalType读取JSON,而是先以DoubleType读取数据,再通过类型转换得到DecimalType,具体步骤如下:
- 读取JSON时使用默认推断的schema(或显式用DoubleType定义
C_0_1) - 通过
cast方法将C_0_1转换为DecimalType(10,6) - 写入Parquet时无需额外指定schema
修正后的代码:
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, DoubleType, DecimalType, ArrayType # 初始化SparkSession(替代旧的SparkContext/HiveContext,Spark 2.0+标准API) spark = SparkSession.builder.appName("JsonToParquet").getOrCreate() json_str = '{"C_0_0": "INDIA", "C_0_1": 1219.765432, "C_0_2": {"C_1_0": "INDIA", "C_1_1": ["INDIA", "INDIA"]}}' rdd = spark.sparkContext.parallelize([json_str]) # 1. 先读取为带DoubleType的DataFrame nested_df = spark.read.json(rdd) print("Original Schema: ", nested_df.schema) # 2. 将C_0_1转换为DecimalType(10,6) converted_df = nested_df.withColumn("C_0_1", nested_df["C_0_1"].cast(DecimalType(10,6))) print("Converted Schema: ", converted_df.schema) # 3. 写入Parquet,无需额外指定schema file_location = "path/to/your/parquet" converted_df.repartition(1).write.parquet(file_location, mode="overwrite") # 读取验证 parDf = spark.read.parquet(file_location) parDf.show(20, False) parDf.printSchema()
额外说明
- 优先使用
SparkSession替代旧的SparkContext/HiveContext/SQLContext,这是Spark 2.0及以上版本的标准API,兼容性和稳定性更好。 - 如果必须在读取阶段就指定DecimalType,需要先将JSON中的数值以字符串形式存储,再用DecimalType读取,但这种方式需要修改原始JSON结构,灵活性不如转换方式。
内容的提问来源于stack exchange,提问作者Abhik NASKAR
相关产品推荐
相关产品推荐

