You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 05:00:57