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表达式时参数不匹配,通常由以下原因导致:
- Spark与Scala版本不兼容:Lambda的序列化逻辑依赖Scala和JDK版本,版本不匹配会导致反序列化时参数数量异常。
- 隐式重载引发的Lambda生成问题:直接传递单个字符串给
dropDuplicates时,Scala的隐式转换可能生成不符合Spark预期的Lambda实例。 - 类继承结构中的序列化缺陷:当前类重写了
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
相关产品推荐
相关产品推荐

