Spark 2.4.0 PySpark写入Parquet触发ArrayIndexOutOfBoundsException排查
问题排查与解决方案
可能的原因分析
- Spark 2.4.0
sequence函数的Date类型兼容性bug:Spark 2.4.x版本中,sequence函数结合DateType和interval 1 month生成的序列,在写入Parquet时会触发底层序列化异常。虽然show()能正常展示数据,但写入操作的序列化逻辑会暴露数组越界问题。 - Decimal类型列的Parquet写入冲突:
SK_APPLICATION为Decimal类型,与sequence生成的DATE列(可能被隐式转为TimestampType)组合写入时,Parquet writer的类型处理逻辑出现异常。
验证步骤
先确认missing_dates的Schema是否符合预期:
missing_dates.printSchema()
正常输出应为:
root |-- SK_APPLICATION: decimal(...) (nullable = true) |-- DATE: date (nullable = true) |-- CODE: string (nullable = false)
如果DATE列是timestamp类型,说明sequence函数发生了隐式类型转换,这大概率是问题根源。
解决方案
方案1:显式转换DATE列为DateType
将sequence的输入转为TimestampType,生成序列后再转回DateType,规避隐式类型转换的问题:
from pyspark.sql.types import DateType, TimestampType missing_dates = (grouped_df .select("SK_APPLICATION", explode(sequence( col("min_date").cast(TimestampType()), col("max_date").cast(TimestampType()), expr("interval 1 month") )).cast(DateType()).alias("DATE")) .withColumn("CODE", lit("#")) ) missing_dates.write.mode('overwrite').parquet('missing_dates')
方案2:转换SK_APPLICATION为Long类型
Decimal类型在Parquet写入时可能存在兼容性问题,将其转为Long类型后再执行写入:
from pyspark.sql.types import LongType missing_dates = (grouped_df .withColumn("SK_APPLICATION", col("SK_APPLICATION").cast(LongType())) .select("SK_APPLICATION", explode(sequence(col("min_date"), col("max_date"), expr("interval 1 month"))).alias("DATE")) .withColumn("CODE", lit("#")) ) missing_dates.write.mode('overwrite').parquet('missing_dates')
方案3:替换sequence函数,用add_months生成序列
如果上述方案无效,可放弃sequence,改用递归方式生成月度序列(适配Spark 2.4.x版本):
from pyspark.sql.types import IntegerType def generate_monthly_dates(grouped_df): # 计算每个SK_APPLICATION的月份差 df = grouped_df.withColumn("months_diff", months_between(col("max_date"), col("min_date")).cast(IntegerType())) # 生成月份偏移量序列,逐个计算对应日期 df = df.withColumn("month_offset", explode(sequence(lit(0), col("months_diff")))) df = df.withColumn("DATE", add_months(col("min_date"), col("month_offset"))) return df.select("SK_APPLICATION", "DATE").withColumn("CODE", lit("#")) missing_dates = generate_monthly_dates(grouped_df) missing_dates.write.mode('overwrite').parquet('missing_dates')
额外检查
提前过滤grouped_df中min_date > max_date的异常行,避免生成无效序列:
grouped_df = grouped_df.filter(col("min_date") <= col("max_date"))
内容的提问来源于stack exchange,提问作者Владислав Черкасов
相关产品推荐
相关产品推荐

