读取S3上Parquet文件并跳过不符合Schema的行的方法
解决Parquet文件加载时类型不匹配及过滤无效行的问题
问题背景
存储在S3上的一批Parquet文件(由gzip压缩CSV生成),需要加载为Spark DataFrame,同时需跳过以下无效行:
- col_1为空或无法转换为Long类型(原数据为字符串或其他非数值类型)
- col_2为空
指定Schema如下:
schema = StructType() \ .add("col_1",LongType(),True) \ .add("col_2",StringType(),True) \ .add("col_3",StringType(),True)
使用常规加载语句时触发报错:
Error while reading file XXXXX. Parquet column cannot be converted. Column: [col_1], Expected: LongType, Found: BINARY
Caused by: SchemaColumnConvertNotSupportedException
解决方案
1. 先以兼容模式加载全量数据
直接强制指定目标Schema会因类型不匹配失败,需先以兼容方式加载所有行,避免转换错误中断读取:
from pyspark.sql.types import StructType, LongType, StringType, BinaryType # 方式1:不指定Schema,让Spark自动推断(适合字段结构稳定的情况) temp_df = spark.read.format("parquet").load('s3_bucket/folder_parquet_files/*') # 方式2:显式定义兼容Schema,将col_1设为BinaryType(更可控) # temp_schema = StructType() \ # .add("col_1", BinaryType(), True) \ # .add("col_2", StringType(), True) \ # .add("col_3", StringType(), True) # temp_df = spark.read.format("parquet").schema(temp_schema).load('s3_bucket/folder_parquet_files/*')
2. 转换类型并过滤无效行
使用try_cast安全转换col_1类型,同时过滤不符合要求的行:
from pyspark.sql.functions import col, try_cast, trim # 第一步:过滤col_2为空的行(含空白字符串) filtered_df = temp_df.filter( (trim(col("col_2")).isNotNull()) & (trim(col("col_2")) != "") ) # 第二步:将col_1转换为Long类型,转换失败则返回null converted_df = filtered_df.withColumn( "col_1", try_cast(col("col_1").cast("string"), "long") ) # 第三步:过滤col_1转换后为空的行(原数据为空或无法转为Long的情况) final_df = converted_df.filter(col("col_1").isNotNull())
3. 对齐目标Schema(可选)
如果需要严格匹配最初指定的Schema,可强制转换字段类型:
final_df = final_df.select( col("col_1").cast(LongType()).alias("col_1"), col("col_2").cast(StringType()).alias("col_2"), col("col_3").cast(StringType()).alias("col_3") )
4. 验证结果
查看加载后的数据结构和内容:
final_df.printSchema() final_df.show()
关键细节
try_cast是核心:它会尝试转换类型,失败时返回null而非抛出异常,完美处理类型不兼容的行- 先过滤col_2无效行:减少后续类型转换的计算量,提升效率
- 二进制转字符串再转Long:解决Parquet中col_1以二进制存储字符串的类型冲突问题
内容的提问来源于stack exchange,提问作者Antonius
相关产品推荐
相关产品推荐

