使用parquet-tool读取Parquet文件时字符串值转科学计数法的问题
问题:文本转Parquet后数值格式异常(科学计数法+精度丢失)
将文本文件读取为Spark DataFrame后写入Parquet,使用parquet-tool读取时,原文本中的数值出现科学计数法转换及精度丢失问题:
- 原文本示例数据:
111173245.136459|131856.12 - 写入Parquet核心代码:
df.write.format("parquet").save(test.parquet)
执行parquet-tools show --columns col1,col2 ./test.parquet读取时,Column1显示为1.11073e+08(科学计数法),Column2显示为130850,与原数值完全不符。已为两列指定StringType Schema,且未对列做任何处理,不清楚写入Parquet时自动转换格式的原因,需了解解决方案或需添加的写入参数。
完整示例代码
import io,csv import base64 import pandas as pd import pyspark.sql.types from pyspark.context import SparkContext from pyspark.sql import SparkSession from pyspark.sql.functions import monotonically_increasing_id from pyspark.sql.types import * ################################################### spark = SparkSession.builder.master('local').config('spark.sql.session.timeZone', 'UTC').getOrCreate() sc = spark.sparkContext ################################################### rDelim='\r\n' cDelim='|' rDelim = rDelim.replace('\\n','\n') rDelim = rDelim.replace('\\r','\r') encoding = 'UTF-8' ################################################### fileName='test.bcp' col=['Column1|StringType','Column2|StringType']; ################################################### def typeName(dType): return getattr(pyspark.sql.types, dType)() ################################################### def getType(): cFields=[] for x in col: colList = x.split("|") cFields.append(StructField(colList[0], typeName(colList[1]), True)) schema = StructType(cFields) return schema ################################################### def csvStr(x): output = io.StringIO("") csv.writer(output).writerow(x) return output.getvalue().strip() ################################################### def creataeParquet(schema): rdd = sc.newAPIHadoopFile(fileName, "org.apache.hadoop.mapreduce.lib.input.TextInputFormat","org.apache.hadoop.io.LongWritable", "org.apache.hadoop.io.Text", conf={"textinputformat.record.delimiter": rDelim}).map(lambda line: line[1].split(cDelim)) df = spark.read.format( "com.databricks.spark.csv").schema(schema).option( "escape", '\"').option( "quote", '\"').option( "header", "false").option( "encoding", encoding).csv( rdd.map(csvStr)) df.printSchema(); df.show(2000, truncate=False) df.write.format("parquet").save('/tmp/test.parquet') schema = getType() cnt = creataeParquet(schema)
问题原因及解决方法
原因分析
问题出在数据读取阶段,而非Parquet写入环节:
csvStr函数中使用csv.writer处理分割后的字段时,会自动将数值型字符串转换为浮点数格式输出,导致原始字符串被修改(比如转成科学计数法、丢失精度)。- 即使指定了
StringTypeSchema,Spark CSV读取时可能仍存在隐性类型推断,覆盖了手动指定的类型,进一步加剧格式异常。
解决方法
方法1:绕过csv.writer,手动构造CSV格式字符串
删除csvStr函数,直接手动拼接符合CSV规范的字符串,避免自动类型转换:
# 修改rdd生成逻辑,替换原map(csvStr) rdd = sc.newAPIHadoopFile(fileName, "org.apache.hadoop.mapreduce.lib.input.TextInputFormat","org.apache.hadoop.io.LongWritable", "org.apache.hadoop.io.Text", conf={"textinputformat.record.delimiter": rDelim}) \ .map(lambda line: line[1].split(cDelim)) \ .map(lambda x: ','.join(f'"{item}"' for item in x)) # 手动添加引号与分隔符
方法2:强制禁用Spark CSV的类型推断
在读取CSV时添加inferSchema=false参数,确保严格使用指定的Schema:
df = spark.read.format("com.databricks.spark.csv") \ .schema(schema) \ .option("escape", '\"') \ .option("quote", '\"') \ .option("header", "false") \ .option("encoding", encoding) \ .option("inferSchema", "false") # 新增该参数,禁用自动类型推断 .csv(rdd.map(csvStr))
方法3:简化读取逻辑,跳过中间CSV转换
直接使用Spark读取文本文件并分割字段,避免复杂的Hadoop API处理:
def creataeParquet(schema): # 直接读取文本文件并分割字段 df = spark.read.text(fileName) \ .selectExpr("split(value, '\\|') as cols") \ .select( col("cols")[0].alias("Column1"), col("cols")[1].alias("Column2") ) \ .cast(schema) # 强制转换为指定Schema df.printSchema() df.show(2000, truncate=False) df.write.format("parquet").save('/tmp/test.parquet')
验证步骤
修改代码后,先通过df.show()确认DataFrame中的数据与原文本完全一致,再写入Parquet文件,最后用parquet-tools检查即可看到原始字符串格式,无科学计数法及精度丢失问题。
内容的提问来源于stack exchange,提问作者SwapnilM
相关产品推荐
相关产品推荐

