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
相关产品推荐
相关产品推荐

