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

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异常。

可行解决办法:

  1. 升级Connector版本:升级到3.4.0及以上版本,该问题已在后续版本中修复。新版本针对空列表的schema匹配逻辑增加了判空分支,避免对空集合调用reduce操作。
  2. 临时规避方案:如果暂时无法升级,可以在加载数据前过滤掉包含空列表的记录,或者在数据处理阶段对空列表做转换。例如,在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"))
)
  1. 自定义类型转换器:如果需要保留空列表,可以自定义类型转换器逻辑,覆盖原Connector的TypeConverter类,在处理空列表时直接返回对应的Spark类型(如ArrayType(StringType)),跳过reduce操作。

内容的提问来源于stack exchange,提问作者ClogShug

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 22:32:11