You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.16 05:10:21