Spark中使用System.out作为RDD任务报错:任务不可序列化
嘿,这个问题我之前刚踩过坑!让我给你理清楚原因和可行的解决办法~
问题原因
Spark的转换和行动算子(比如foreach)中用到的对象,需要被序列化后发送到集群的Worker节点上执行。而你代码里用的System.out::println这个方法引用,实际上会捕获System.out这个PrintStream实例——但java.io.PrintStream并没有实现Serializable接口,所以Spark尝试序列化它的时候就抛出了NotSerializableException。
至于你看的教程能正常运行,大概率是教程用的是Scala(Scala的println是独立函数,不会绑定外部的PrintStream实例),或者教程里用了其他避开序列化的写法。
解决办法
这里有两种靠谱的方案,你可以根据场景选择:
方案1:使用foreachPartition替代foreach
foreachPartition是针对每个分区执行操作,在Worker节点的本地进程中直接获取System.out,不需要序列化Driver端的PrintStream对象。修改后的代码如下:
myRDD.foreachPartition(iter -> { // 这里的System.out是Worker节点本地的输出流,无需序列化 while (iter.hasNext()) { System.out.println(iter.next()); } });
这种方法适合大数据量场景,因为不会把数据拉到Driver端,而是直接在Worker节点处理输出。
方案2:将RDD数据收集到Driver端后打印(仅适合小数据测试)
如果只是本地测试、数据量很小,可以先把RDD的数据拉到Driver节点,再执行打印:
// collect()会把RDD所有数据拉到Driver内存,数据量大时慎用 myRDD.collect().forEach(System.out::println);
注意:这种方法只适合小数据集,否则Driver端会因为内存不足崩溃。
额外提示
以后写Spark代码时要记住:算子内部引用的外部对象必须是可序列化的,如果遇到不可序列化的对象,要么换成本地创建(比如foreachPartition里的System.out),要么把对象改成可序列化的类型。
内容的提问来源于stack exchange,提问作者ashender

