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

Spark-Scala关联查询优化:如何对关联结果进行分组?

Hey 你好!看起来你现在遇到的是inner join带来的笛卡尔积问题——因为Data1和Data2里的Code值完全一致,所以关联后生成了所有Id和Date的组合,也就是4×2=8行结果。如果你的目标是减少输出行数,通过分组聚合来整合这些重复关联的信息,这里有几种不同的优化方案,取决于你最终想要的输出格式:

方案1:按Code聚合,把所有Id和Date整合为数组

如果希望同一个Code下的所有Id和Date都放在一行里,这样最终只会输出1行结果(因为Code唯一):

import org.apache.spark.sql.functions.{collect_list, col}

Data1.join(Data2, Seq("Code"), "inner")
  .groupBy("Code")
  .agg(
    collect_list("Id").alias("Ids"),
    collect_list("Date").alias("Dates")
  )
  .select("Code", "Ids", "Dates")
  .show(false)

输出结果示例:

+------------+------------------------+----------------+
|Code        |Ids                     |Dates           |
+------------+------------------------+----------------+
|06869242986|[0839, 4395, 3796, 3592]|[202050, 202051]|
+------------+------------------------+----------------+

方案2:按Id分组,聚合对应的Date列表

如果希望每个Id对应所有匹配的Date,这样最终输出4行(对应Data1里的4个唯一Id):

Data1.join(Data2, Seq("Code"), "inner")
  .groupBy("Id", "Code")
  .agg(collect_list("Date").alias("Dates"))
  .select("Id", "Dates", "Code")
  .show(false)

输出结果示例:

+----+----------------+------------+
|Id  |Dates           |Code        |
+----+----------------+------------+
|0839|[202050, 202051]|06869242986|
|4395|[202050, 202051]|06869242986|
|3796|[202050, 202051]|06869242986|
|3592|[202050, 202051]|06869242986|
+----+----------------+------------+

方案3:按Date分组,聚合对应的Id列表

如果希望每个Date对应所有匹配的Id,这样最终输出2行(对应Data2里的2个唯一Date):

Data1.join(Data2, Seq("Code"), "inner")
  .groupBy("Date", "Code")
  .agg(collect_list("Id").alias("Ids"))
  .select("Date", "Ids", "Code")
  .show(false)

输出结果示例:

+------+------------------------+------------+
|Date  |Ids                     |Code        |
+------+------------------------+------------+
|202050|[0839, 4395, 3796, 3592]|06869242986|
|202051|[0839, 4395, 3796, 3592]|06869242986|
+------+------------------------+------------+

额外小提示:

  • 如果你的Data1/Data2里可能存在重复的Id或Date值,可以用collect_set代替collect_list,这样聚合后的集合会自动去重。
  • 这些方案的核心是用分组聚合替代原始的笛卡尔积式关联结果,不仅减少了输出行数,在数据量较大时也能提升处理效率,避免不必要的数据膨胀。

内容的提问来源于stack exchange,提问作者Rebe

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 08:37:39