Java循环读取Spark Dataset行时无法读取内部其他Dataset如何解决
问题根本原因
ds.foreach() 属于Spark分布式算子,内部的函数逻辑会被序列化后发送到各个Executor节点执行。你在foreach内部读取其他Dataset的操作,本质需要调用Driver端的SparkSession/SparkContext执行算子,但Executor节点上没有可用的SparkContext实例,所以会读取失败。
而ds.toJavaRDD().collect().forEach()的写法中,collect()会先把整个Dataset的数据拉取到Driver节点的内存中,后续的forEach是JDK自带的集合遍历逻辑,完全在Driver端执行,因此可以正常调用SparkSession读取其他Dataset。
可行解决方案
方案1:保留现有collect逻辑(适合小数据集场景)
如果你的外层ds数据量较小,不会撑爆Driver内存,直接保留现有写法即可,是成本最低的方案。
注意:如果
ds数据量超过Driver内存配置会触发OOM,该场景下不能使用该方案。
方案2:提前对关联Dataset做Join(适合大数据集场景)
如果你遍历外层ds每行的目的是拿行内字段去其他Dataset查询关联数据,直接用Spark的Join算子提前关联两个Dataset,不需要写循环逐行查询,完全通过Spark分布式算子实现,性能最优,也不会出现跨节点调用SparkContext的问题。
示例代码结构:
// 假设外层ds为userDs,要关联的其他Dataset为orderDs,关联键为user_id Dataset<Row> joinedDs = userDs.join(orderDs, userDs.col("user_id").equalTo(orderDs.col("user_id")), "left"); // 直接对关联后的joinedDs做后续处理,不需要循环逐行查询 joinedDs.foreach(row -> { // 此处已能获取两个表的所有字段,无需再读取其他Dataset // 处理逻辑仅可使用行内已有的数据,不能调用Spark相关API });
方案3:将小维度Dataset转为广播变量(适合小维度表场景)
如果要读取的其他Dataset是小尺寸的维度表,数据量不大,可以提前将该Dataset的数据collect后广播到所有Executor节点,在ds.foreach()中直接读取广播变量内的数据,不需要在Executor端调用SparkSession读取Dataset。
示例代码结构:
// 提前将小维度表collect后转为广播变量 List<Row> dimList = dimDs.collectAsList(); Broadcast<List<Row>> dimBroadcast = javaSparkContext.broadcast(dimList); // 遍历外层大Dataset时直接读取广播变量 ds.foreach(row -> { List<Row> dimData = dimBroadcast.value(); // 从dimData中查询所需的对应数据即可 });
禁止操作
绝对不要尝试在ds.foreach()、ds.map()这类分布式算子内部创建SparkSession、调用Dataset相关API,这类操作仅可在Driver端执行,放到Executor中一定会失败。
内容的提问来源于stack exchange,提问作者Nikhil Agrawal

