Scala中使用.join实现两个Dataset按ID关联取匹配数据
Scala Spark 基于Dataset.join实现ID交集取左表记录方案
该方案等价于指定的内连接SQL逻辑,执行效率远高于isin类实现,适配万级以上数据量的处理场景。
前置定义
使用强类型Dataset需要先定义对应数据结构的样例类,关联阶段仅用到Dataset2的id字段,可提前裁剪冗余字段减少计算开销:
// Dataset1对应结构 case class Ds1Record(id: Long, ownerId: String, marketplace: Int) // Dataset2对应结构 case class Ds2Record(id: Long, marketplace: Int)
核心实现代码
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions.broadcast val spark = SparkSession.builder().appName("IdIntersectionJoin").getOrCreate() import spark.implicits._ // 加载/构造两份数据集(示例为直接造数,实际场景替换为读文件/表的逻辑即可) val ds1 = Seq( Ds1Record(1234, "george", 1), Ds1Record(2345, "mike", 1), Ds1Record(3456, "anish", 1), Ds1Record(4567, "annie", 1), Ds1Record(5678, "waker", 2) ).toDS() val ds2 = Seq( Ds2Record(1234, 1), Ds2Record(2345, 1), Ds2Record(8888, 1), Ds2Record(9999, 1), Ds2Record(7777, 1) ).toDS() // 核心join逻辑:内连接+关联键id+提前裁剪右表冗余字段 val resultDs = ds1.join( right = ds2.select("id"), // 仅保留右表关联需要的id列,降低shuffle开销 usingColumns = Seq("id"), // 指定关联键,自动合并同名字段避免列冲突 joinType = "inner" // 内连接天然保留两边id同时存在的记录 ).as[Ds1Record] // 转回强类型Dataset
结果验证
执行resultDs.show()即可得到符合需求的输出:
+----+-------+-----------+ | id|ownerId|marketplace| +----+-------+-----------+ |1234| george| 1| |2345| mike| 1| +----+-------+-----------+
关键说明
- 关联前裁剪右表非必要字段是通用优化手段,shuffle阶段传输的数据量越小,执行速度越快
- 用
Seq("id")指定关联键的写法,比手动写ds1("id") === ds2("id")的条件写法更简洁,不会出现结果集中存在重复id列的问题 - 性能特性:该实现底层会自动根据表大小选择BroadcastJoin或SortMergeJoin执行计划,万级到亿级数据量都能稳定运行,完全避免
isin方案需要将右表全量id收集到Driver端带来的OOM风险 - 如果Dataset2数据量远小于Dataset1(通常阈值为10万条以内),可给右表加
broadcast()标记触发广播连接,跳过shuffle阶段进一步提升执行速度:val optimizedResult = ds1.join(broadcast(ds2.select("id")), Seq("id"), "inner").as[Ds1Record]
内容的提问来源于stack exchange,提问作者Jericho Sims
相关产品推荐
相关产品推荐

