Aerospike Spark Connector加载数据失败:空列表引发reduceLeft异常求助
问题描述
我使用Aerospike Spark Connector从Aerospike集群加载数据到DataFrame,处理后写入另一个Aerospike集群。数据包含两个bin:一个是字符串类型的列表,另一个是键为字符串、值为长整型的Map。运行Spark应用时任务失败,驱动端堆栈跟踪如下:
24/06/21 09:59:58 WARN TaskSetManager: Lost task 0.0 in stage 0.0 (TID 0) (10.10.16.151 executor 0): java.lang.UnsupportedOperationException: empty.reduceLeft at scala.collection.TraversableOnce.reduceLeft(TraversableOnce.scala:185) at scala.collection.TraversableOnce.reduceLeft$(TraversableOnce.scala:183) at scala.collection.mutable.ArrayBuffer.scala$collection$IndexedSeqOptimized$$super$reduceLeft(ArrayBuffer.scala:49) at scala.collection.IndexedSeqOptimized.reduceLeft(IndexedSeqOptimized.scala:77) at scala.collection.IndexedSeqOptimized.reduceLeft$(IndexedSeqOptimized.scala:76) at scala.collection.mutable.ArrayBuffer.reduceLeft(ArrayBuffer.scala:49) at scala.collection.TraversableOnce.reduce(TraversableOnce.scala:213) at scala.collection.TraversableOnce.reduce$(TraversableOnce.scala:213) at scala.collection.AbstractTraversable.reduce(Traversable.scala:108) at com.aerospike.spark.converters.TypeConverter$.matchesSchemaType(TypeConverter.scala:246) at com.aerospike.spark.converters.TypeConverter$.convertToSparkType(TypeConverter.scala:353) at com.aerospike.spark.converters.TypeConverter$.binToValue(TypeConverter.scala:428) at com.aerospike.spark.sql.sources.v2.RowIterator.$anonfun$get$2(RowIterator.scala:60) at scala.collection.TraversableLike.$anonfun$map$1(TraversableLike.scala:238) at scala.collection.IndexedSeqOptimized.foreach(IndexedSeqOptimized.scala:36) at scala.collection.IndexedSeqOptimized.foreach$(IndexedSeqOptimized.scala:33) at scala.collection.mutable.ArrayOps$ofRef.foreach(ArrayOps.scala:198) at scala.collection.TraversableLike.map(TraversableLike.scala:238) at scala.collection.TraversableLike.map$(TraversableLike.scala:231) at scala.collection.mutable.ArrayOps$ofRef.map(ArrayOps.scala:198) at com.aerospike.spark.sql.sources.v2.RowIterator.get(RowIterator.scala:48) at com.aerospike.spark.sql.sources.v2.RowIterator.get(RowIterator.scala:21) at org.apache.spark.sql.execution.datasources.v2.PartitionIterator.next(DataSourceRDD.scala:89) at org.apache.spark.sql.execution.datasources.v2.MetricsRowIterator.next(DataSourceRDD.scala:124) at org.apache.spark.sql.execution.datasources.v2.MetricsRowIterator.next(DataSourceRDD.scala:121) at org.apache.spark.InterruptibleIterator.next(InterruptibleIterator.scala:40) at scala.collection.Iterator$$anon$10.next(Iterator.scala:459) 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(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:750)
排查结果
该异常是由于对值为空列表的bin执行schema匹配时调用reduce函数导致的。
版本信息
- Spark版本:3.1.2
- Connector版本:3.3.0
疑问
是否存在解决办法,或者这是Spark Connector的自身问题?
解决方案
这是Aerospike Spark Connector 3.3.0版本的已知bug,根源在于TypeConverter.scala的matchesSchemaType方法中,对空列表执行reduce操作时没有做判空处理,导致抛出empty.reduceLeft异常。
可行解决办法:
- 升级Connector版本:升级到3.4.0及以上版本,该问题已在后续版本中修复。新版本针对空列表的schema匹配逻辑增加了判空分支,避免对空集合调用reduce操作。
- 临时规避方案:如果暂时无法升级,可以在加载数据前过滤掉包含空列表的记录,或者在数据处理阶段对空列表做转换。例如,在Spark DataFrame中将空列表替换为包含默认值的列表:
import org.apache.spark.sql.functions._ // 根据业务需求处理空列表,此处示例替换为含默认字符串的列表 val processedDF = originalDF.withColumn( "list_bin", when(size(col("list_bin")) === 0, array(lit("default_value"))).otherwise(col("list_bin")) )
- 自定义类型转换器:如果需要保留空列表,可以自定义类型转换器逻辑,覆盖原Connector的
TypeConverter类,在处理空列表时直接返回对应的Spark类型(如ArrayType(StringType)),跳过reduce操作。
内容的提问来源于stack exchange,提问作者ClogShug
相关产品推荐
相关产品推荐

