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

如何使用Spark DSL实现无重复自交叉连接并输出目标结果

Spark DSL 实现无重复国家两两配对方案

前置假设

你已初始化Spark运行环境,且存在存储国家取值的基础DataFrame,结构如下:

root
 |-- country: string (nullable = true)

包含的取值为:IN、PK、AU、SL。

核心实现逻辑

通过自连接时添加country_a < country_b的不等值过滤条件,一次性实现两个效果:

  • 排除国家自身和自身配对的无效结果
  • 避免出现顺序相反的重复配对(仅保留字典序更小的国家在前的组合,不会同时出现IN vs PK和PK vs IN)

完整实现代码

PySpark 版本

from pyspark.sql.functions import col, concat, lit

# 步骤1:为原DF分别重命名列,用于自连接区分左右表
df_left = df.select(col("country").alias("country_a"))
df_right = df.select(col("country").alias("country_b"))

# 步骤2:自连接+核心过滤条件
result_df = df_left.join(
    df_right,
    on=col("country_a") < col("country_b"),
    how="inner"
)

# 步骤3:拼接为指定的「X vs Y」格式输出
result_df = result_df.select(
    concat(col("country_a"), lit(" vs "), col("country_b")).alias("match_result")
)

# 验证输出
result_df.show(truncate=False)

Scala Spark 版本

import org.apache.spark.sql.functions.{col, concat, lit}

// 步骤1:为原DF分别重命名列,用于自连接区分左右表
val dfLeft = df.select(col("country").as("country_a"))
val dfRight = df.select(col("country").as("country_b"))

// 步骤2:自连接+核心过滤条件
val resultDf = dfLeft.join(
  dfRight,
  col("country_a") < col("country_b"),
  "inner"
)

// 步骤3:拼接为指定的「X vs Y」格式输出
val finalDf = resultDf.select(
  concat(col("country_a"), lit(" vs "), col("country_b")).as("match_result")
)

// 验证输出
finalDf.show(false)

输出结果示例

+------------+
|match_result|
+------------+
|AU vs IN    |
|AU vs PK    |
|AU vs SL    |
|IN vs PK    |
|IN vs SL    |
|PK vs SL    |
+------------+

本实现逻辑和你参考的SQL实现效果完全等价。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 20:45:05