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

Spark使用dropDuplicates函数时出现序列化错误求助

Spark dropDuplicates 序列化异常排查与解决

问题描述

在Scala代码中调用Spark的dropDuplicates函数时遇到序列化异常,代码如下:

override def innerTransform(dataFrames: Map[ReaderKey, DataFrame]): DataFrame = {
    val fromKafka = super.innerTransform(dataFrames)
    fromKafka.dropDuplicates("column8")
}

报错堆栈

[java.io.InvalidObjectException: ReflectiveOperationException during deserialization
    at java.base/java.lang.invoke.SerializedLambda.readResolve(SerializedLambda.java:280)
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77)
    at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.base/java.lang.reflect.Method.invoke(Method.java:568)
    at java.base/java.io.ObjectStreamClass.invokeReadResolve(ObjectStreamClass.java:1190)
    at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2266)
    at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733)
    at java.base/java.io.ObjectInputStream.readArray(ObjectInputStream.java:2157)
    at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1721)
    at java.base/java.io.ObjectInputStream$FieldValues.<init>(ObjectInputStream.java:2606)
    at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2457)
    at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257)
    at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733)
    at java.base/java.io.ObjectInputStream$FieldValues.<init>(ObjectInputStream.java:2606)
    at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2457)
    at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257)
    at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733)
    at java.base/java.io.ObjectInputStream$FieldValues.<init>(ObjectInputStream.java:2606)
    at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2457)
    at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257)
    at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733)
    at java.base/java.io.ObjectInputStream.readObject(ObjectInputStream.java:509)
    at java.base/java.io.ObjectInputStream.readObject(ObjectInputStream.java:467)
    at scala.collection.immutable.List$SerializationProxy.readObject(List.scala:488)
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77)
    at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.base/java.lang.reflect.Method.invoke(Method.java:568)
    at java.base/java.io.ObjectStreamClass.invokeReadObject(ObjectStreamClass.java:1100)
    at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2423)
    at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257)
    at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733)
    at java.base/java.io.ObjectInputStream$FieldValues.<init>(ObjectInputStream.java:2606)
    at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2457)
    at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257)
    at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733)
    at java.base/java.io.ObjectInputStream$FieldValues.<init>(ObjectInputStream.java:2606)
    at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2457)
    at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257)
    at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733)
    at java.base/java.io.ObjectInputStream.readObject(ObjectInputStream.java:509)
    at java.base/java.io.ObjectInputStream.readObject(ObjectInputStream.java:467)
    at scala.collection.immutable.List$SerializationProxy.readObject(List.scala:488)
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77)
    at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.base/java.lang.reflect.Method.invoke(Method.java:568)
    at java.base/java.io.ObjectStreamClass.invokeReadObject(ObjectStreamClass.java:1100)
    at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2423)
    at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257)
    at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733)
    at java.base/java.io.ObjectInputStream$FieldValues.<init>(ObjectInputStream.java:2606)
    at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2457)
    at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257)
    at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733)
    at java.base/java.io.ObjectInputStream$FieldValues.<init>(ObjectInputStream.java:2606)
    at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2457)
    at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257)
    at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733)
    at java.base/java.io.ObjectInputStream.readObject(ObjectInputStream.java:509)
    at java.base/java.io.ObjectInputStream.readObject(ObjectInputStream.java:467)
    at scala.collection.immutable.List$SerializationProxy.readObject(List.scala:488)
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77)
    at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.base/java.lang.reflect.Method.invoke(Method.java:568)
    at java.base/java.io.ObjectStreamClass.invokeReadObject(ObjectStreamClass.java:1100)
    at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2423)
    at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257)
    at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733)
    at java.base/java.io.ObjectInputStream$FieldValues.<init>(ObjectInputStream.java:2606)
    at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2457)
    at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257)
    at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733)
    at java.base/java.io.ObjectInputStream$FieldValues.<init>(ObjectInputStream.java:2606)
    at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2457)
    at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257)
    at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733)
    at java.base/java.io.ObjectInputStream.readObject(ObjectInputStream.java:509)
    at java.base/java.io.ObjectInputStream.readObject(ObjectInputStream.java:467)
    at org.apache.spark.serializer.JavaDeserializationStream.readObject(JavaSerializer.scala:75)
    at org.apache.spark.serializer.JavaSerializerInstance.deserialize(JavaSerializer.scala:114)
    at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:88)
    at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:55)
    at org.apache.spark.scheduler.Task.run(Task.scala:123)
    at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:411)
    at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:414)
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
    at java.base/java.lang.Thread.run(Thread.java:840)
Caused by: java.lang.reflect.InvocationTargetException
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77)
    at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.base/java.lang.reflect.Method.invoke(Method.java:568)
    at java.base/java.lang.invoke.SerializedLambda.readResolve(SerializedLambda.java:278)
    ... 91 more
Caused by: java.lang.IllegalArgumentException: too many arguments
    at java.base/java.lang.invoke.LambdaMetafactory.altMetafactory(LambdaMetafactory.java:511)
    at scala.runtime.LambdaDeserializer$.makeCallSite$1(LambdaDeserializer.scala:105)
    at scala.runtime.LambdaDeserializer$.deserializeLambda(LambdaDeserializer.scala:114)
    at scala.runtime.LambdaDeserialize.deserializeLambda(LambdaDeserialize.java:38)
    at org.apache.spark.sql.execution.LocalTableScanExec.$deserializeLambda$(LocalTableScanExec.scala)
    ... 96 more]

DataFrame Schema

root
 |-- column1: struct (nullable = true)
 |    |-- column1: binary (nullable = true)
 |    |-- column2: array (nullable = true)
 |    |    |-- element: struct (containsNull = true)
 |    |    |    |-- column1: string (nullable = true)
 |    |    |    |-- column2: binary (nullable = true)
 |    |-- column3: string (nullable = true)
 |    |-- column4: integer (nullable = false)
 |    |-- column5: long (nullable = false)
 |    |-- column6: timestamp (nullable = true)
 |    |-- column7: integer (nullable = false)
 |-- column8: string (nullable = true)
 |-- column9: string (nullable = true)
 |-- column10: string (nullable = true)
 |-- column11: string (nullable = true)
 |-- column12: string (nullable = true)
 |-- column13: string (nullable = true)
 |-- column14: string (nullable = true)
 |-- column15: string (nullable = true)
 |-- column16: string (nullable = true)
 |-- column17: integer (nullable = true)
 |-- column18: string (nullable = true)
 |-- column19: string (nullable = true)
 |-- column20: string (nullable = true)
 |-- column21: struct (nullable = true)
 |    |-- column1: string (nullable = true)
 |    |-- column2: string (nullable = true)
 |    |-- column3: string (nullable = true)
 |    |-- column4: string (nullable = true)
 |    |-- column5: string (nullable = true)
 |    |-- column6: string (nullable = true)
 |    |-- column7: string (nullable = true)
 |-- column22: struct (nullable = true)
 |    |-- column1: long (nullable = false)
 |    |-- column2: long (nullable = false)

问题分析

核心报错java.lang.IllegalArgumentException: too many arguments出现在Lambda反序列化阶段,关联到LocalTableScanExec,说明Spark在反序列化执行计划中的Lambda表达式时参数不匹配,通常由以下原因导致:

  1. Spark与Scala版本不兼容:Lambda的序列化逻辑依赖Scala和JDK版本,版本不匹配会导致反序列化时参数数量异常。
  2. 隐式重载引发的Lambda生成问题:直接传递单个字符串给dropDuplicates时,Scala的隐式转换可能生成不符合Spark预期的Lambda实例。
  3. 类继承结构中的序列化缺陷:当前类重写了innerTransform,若父类或当前类存在不可序列化的成员,在dropDuplicates触发的shuffle阶段会引发序列化异常。

解决方案

1. 显式指定列名序列

将dropDuplicates的参数改为Seq包装的列名,避免隐式转换问题:

override def innerTransform(dataFrames: Map[ReaderKey, DataFrame]): DataFrame = {
    val fromKafka = super.innerTransform(dataFrames)
    fromKafka.dropDuplicates(Seq("column8"))
}

2. 确认版本兼容性

检查Spark与Scala版本是否匹配:

  • Spark 2.x:对应Scala 2.11
  • Spark 3.0~3.1:对应Scala 2.12
  • Spark 3.2+:支持Scala 2.12/2.13
    确保编译、打包、运行环境的版本完全一致。

3. 确保类的可序列化性

当前类及其父类必须实现Serializable接口,若存在自定义成员变量,也需保证这些变量可序列化:

class YourTransformClass extends ParentClass with Serializable {
    override def innerTransform(dataFrames: Map[ReaderKey, DataFrame]): DataFrame = {
        val fromKafka = super.innerTransform(dataFrames)
        fromKafka.dropDuplicates(Seq("column8"))
    }
}

4. 简化执行计划(可选)

若上述方法无效,可尝试缓存DataFrame后再执行去重,减少执行计划序列化的复杂度:

override def innerTransform(dataFrames: Map[ReaderKey, DataFrame]): DataFrame = {
    val fromKafka = super.innerTransform(dataFrames).cache()
    val result = fromKafka.dropDuplicates(Seq("column8"))
    result
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 02:25:54