Scala Spark Dataset使用sequence函数触发ArrayIndexOutOfBoundsException问题排查
Spark Dataset生成日期序列时ArrayIndexOutOfBoundsException问题解决
问题背景
我用Scala操作Spark Dataset,需要生成两个格式为yyyy-MM-dd 00:00:000的日期之间按天间隔的序列,使用了以下代码:
Dataset<Row> currentLetters = currentLettersDataset .withColumn("range_dates", functions.sequence(functions.col("start_date"), functions.col("end_date")));
但针对部分日期,执行时会抛出java.lang.ArrayIndexOutOfBoundsException,异常中的索引值与日期间隔的天数相关,堆栈信息如下:
java.lang.ArrayIndexOutOfBoundsException: 135 at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown Source) at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43) at org.apache.spark.sql.execution.WholeStageCodegenExec$$anon$1.hasNext(WholeStageCodegenExec.scala:755) at org.apache.spark.sql.execution.SparkPlan.$anonfun$getByteArrayRdd$1(SparkPlan.scala:345) at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2(RDD.scala:898) at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2$adapted(RDD.scala:898) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:373) at org.apache.spark.rdd.RDD.iterator(RDD.scala:337) at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90) at org.apache.spark.scheduler.Task.run(Task.scala:131) at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:497) at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1439) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:500) at java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source) at java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source) at java.lang.Thread.run(Unknown Source)
问题原因
这个异常是由于Spark内置的sequence函数在代码生成阶段,会预先初始化一个固定大小的数组来存储序列元素。当日期间隔的天数(比如这里的135天)超过了Spark内部预设的数组大小阈值时,就会触发数组越界。
解决方案
方案1:手动生成日期序列(推荐)
放弃使用内置sequence,改用日期差计算+explode的方式生成序列,稳定性更高:
Scala版本
// 1. 计算日期间隔天数 val withDaysDiff = currentLettersDataset.withColumn("days_diff", datediff(col("end_date"), col("start_date"))) // 2. 生成偏移量序列并转换为日期,最后聚合为列表 val result = withDaysDiff .withColumn("day_offset", explode(sequence(lit(0), col("days_diff")))) .withColumn("range_date", date_add(col("start_date"), col("day_offset"))) // 替换为你的原表主键列,保留原有数据结构 .groupBy(col("start_date"), col("end_date")) .agg(collect_list("range_date").alias("range_dates"))
Java版本
// 1. 计算日期间隔天数 Dataset<Row> withDaysDiff = currentLettersDataset .withColumn("days_diff", functions.datediff(functions.col("end_date"), functions.col("start_date"))); // 2. 生成偏移量序列并转换为日期,最后聚合为列表 Dataset<Row> result = withDaysDiff .withColumn("day_offset", functions.explode(functions.sequence(functions.lit(0), functions.col("days_diff")))) .withColumn("range_date", functions.date_add(functions.col("start_date"), functions.col("day_offset"))) // 替换为你的原表主键列,保留原有数据结构 .groupBy(functions.col("start_date"), functions.col("end_date")) .agg(functions.collect_list("range_date").alias("range_dates"));
方案2:调整Spark代码生成阈值
如果必须使用内置sequence,可以通过配置增大代码生成阶段的数组大小阈值:
// 设置阈值为大于你的最大日期间隔天数,例如200 spark.conf.set("spark.sql.codegen.maxFields", "200")
注意:该配置是全局生效的,过大的阈值可能会影响代码生成性能,需根据实际场景权衡。
额外注意事项
- 确保
start_date和end_date是DateType或TimestampType,如果是字符串类型,先转换:
functions.to_timestamp(functions.col("start_date"), "yyyy-MM-dd HH:mm:ssSSS")
- 过滤掉
start_date > end_date的异常数据,避免生成无效序列:
Dataset<Row> filtered = currentLettersDataset.filter(functions.col("start_date").leq(functions.col("end_date")));
内容的提问来源于stack exchange,提问作者Juan Daniel Cortez Rojas
相关产品推荐
相关产品推荐

