为什么Java命令式写法运行正常,等效Lambda表达式在Spark中却报错?
问题核心原因
两个代码段生成的函数对象捕获的外部内容存在本质差异,Spark要求所有传递给算子的函数必须支持Java序列化,报错就是因为代码段2的Lambda实例不满足序列化要求。
前置规则说明
- Spark的分布式执行模型要求,传递给
foreach等算子的函数对象,需要先序列化后分发到各Executor节点执行。Spark提供的ForeachFunction接口默认继承了Serializable接口。 - Java的Lambda/方法引用如果需要序列化,需要同时满足两个条件:
- 转换的目标函数接口实现了
Serializable - Lambda/方法引用捕获的所有外部实例、变量都支持序列化
- 转换的目标函数接口实现了
代码段1正常运行的原因
匿名内部类实现ForeachFunction时,call方法内调用的System.out.println(person)是直接访问System类的静态字段out,匿名内部类不会将System.out作为自身的成员变量捕获,整个匿名内部类实例没有不可序列化的成员,因此可以正常完成序列化分发。
代码段2报错的原因
System.out::println属于绑定实例的方法引用,Java生成对应Lambda实例时,会将用到的System.out(java.io.PrintStream类型)作为捕获成员保存在Lambda实例中。而PrintStream本身没有实现Serializable接口,Spark尝试序列化该Lambda实例时,就会检测到不可序列化的成员,抛出NotSerializableException。
内容的提问来源于stack exchange,提问作者Sheldon Wei
相关产品推荐
相关产品推荐

