Spark读取GCS多Parquet文件时airport_fee类型不一致写入BigQuery失败
解决Spark读取多Parquet文件时列类型不一致的问题
问题根源
Parquet文件的字典编码与数据类型强绑定,当airport_fee列在不同文件中分别以整数、浮点数的字典格式存储时,Spark合并Schema或批量读取时无法兼容两种字典类型,直接触发java.lang.UnsupportedOperationException错误。
可行解决方案
1. 自定义Schema+关闭字典编码
直接定义包含airport_fee为DoubleType的完整Schema,同时关闭Parquet的字典编码优化,强制Spark按指定Schema解析文件,规避类型冲突。
代码示例:
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, DoubleType, StringType, TimestampType, IntegerType import pyspark.sql.functions as F # 定义数据集的完整Schema,将airport_fee明确设为DoubleType custom_schema = StructType([ StructField("VendorID", IntegerType(), True), StructField("tpep_pickup_datetime", TimestampType(), True), StructField("tpep_dropoff_datetime", TimestampType(), True), StructField("passenger_count", DoubleType(), True), StructField("trip_distance", DoubleType(), True), StructField("RatecodeID", DoubleType(), True), StructField("store_and_fwd_flag", StringType(), True), StructField("PULocationID", IntegerType(), True), StructField("DOLocationID", IntegerType(), True), StructField("payment_type", IntegerType(), True), StructField("fare_amount", DoubleType(), True), StructField("extra", DoubleType(), True), StructField("mta_tax", DoubleType(), True), StructField("tip_amount", DoubleType(), True), StructField("tolls_amount", DoubleType(), True), StructField("improvement_surcharge", DoubleType(), True), StructField("total_amount", DoubleType(), True), StructField("congestion_surcharge", DoubleType(), True), StructField("airport_fee", DoubleType(), True) ]) yellow_source = f"gs://{gcp_bucket}/yellow_trip_data/*" spark = SparkSession \ .builder \ .master('yarn') \ .config("spark.sql.files.ignoreCorruptFiles", "true") \ .config("spark.sql.ansi.enabled", "true") \ # 关闭Parquet字典编码,避免类型冲突 .config("spark.sql.parquet.enableDictionary", "false") \ .appName('ny_taxi') \ .getOrCreate() # 读取时指定自定义Schema df = spark.read.schema(custom_schema).parquet(yellow_source) # 后续写入BigQuery df.write \ .mode("overwrite") \ .option("overwriteSchema", "true") \ .format("bigquery") \ .option("temporaryGcsBucket", gcs_spark_bucket) \ .option("dataset", staging_dataset) \ .save("bqtb_stg_yellow")
2. 单文件处理后合并
遍历所有Parquet文件,单独读取并统一airport_fee的类型,再合并为单个DataFrame,绕过批量读取时的字典冲突。
代码示例:
from pyspark.sql import SparkSession import pyspark.sql.functions as F import glob yellow_source_pattern = f"gs://{gcp_bucket}/yellow_trip_data/*.parquet" # 获取所有文件路径 file_paths = glob.glob(yellow_source_pattern) spark = SparkSession \ .builder \ .master('yarn') \ .config("spark.sql.files.ignoreCorruptFiles", "true") \ .config("spark.sql.ansi.enabled", "true") \ .appName('ny_taxi') \ .getOrCreate() merged_df = None for file in file_paths: # 单独读取单个文件 single_df = spark.read.parquet(file) # 强制转换airport_fee为double类型 single_df = single_df.withColumn("airport_fee", F.col("airport_fee").cast("double")) # 合并DataFrame if merged_df is None: merged_df = single_df else: merged_df = merged_df.unionByName(single_df, allowMissingColumns=True) # 写入BigQuery merged_df.write \ .mode("overwrite") \ .option("overwriteSchema", "true") \ .format("bigquery") \ .option("temporaryGcsBucket", gcs_spark_bucket) \ .option("dataset", staging_dataset) \ .save("bqtb_stg_yellow")
3. mergeSchema配合关闭字典编码
开启mergeSchema的同时关闭字典编码,让Spark自动合并Schema后再强制转换类型。
代码示例:
from pyspark.sql import SparkSession import pyspark.sql.functions as F yellow_source = f"gs://{gcp_bucket}/yellow_trip_data/*" spark = SparkSession \ .builder \ .master('yarn') \ .config("spark.sql.files.ignoreCorruptFiles", "true") \ .config("spark.sql.ansi.enabled", "true") \ .config("spark.sql.parquet.enableDictionary", "false") \ .appName('ny_taxi') \ .getOrCreate() # 读取时开启mergeSchema df = spark.read.option("mergeSchema", "true").parquet(yellow_source) # 强制转换airport_fee为double df = df.withColumn("airport_fee", F.col("airport_fee").cast("double")) # 写入BigQuery df.write \ .mode("overwrite") \ .option("overwriteSchema", "true") \ .format("bigquery") \ .option("temporaryGcsBucket", gcs_spark_bucket) \ .option("dataset", staging_dataset) \ .save("bqtb_stg_yellow")
关键提示
spark.sql.parquet.enableDictionary设为false是解决字典类型冲突的核心,关闭后Spark会使用非字典编码读取数据,避免类型不兼容问题。- 自定义Schema时需确保所有列的类型与数据集匹配,避免遗漏或错误定义。
- Yarn集群运行时,需确认配置参数已正确传递至所有Executor节点。
内容的提问来源于stack exchange,提问作者Islam Elsayed
相关产品推荐
相关产品推荐

