PySpark DataFrame同账号内基于时间戳的id配对自连接需求
解决方案:Spark DataFrame生成同账户下时间有序的ID对
这是个典型的同分组内的有序配对问题,咱们可以通过自连接+日期过滤的方式来实现,逻辑清晰还容易调试,下面一步步来:
1. 准备工作:转换日期格式
首先得把字符串类型的time转成Spark的日期类型,这样才能正确比较时间先后:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ // 初始化SparkSession val spark = SparkSession.builder().appName("AccountIdPairs").master("local[*]").getOrCreate() // 构建原始DataFrame val rawDF = spark.createDataFrame(Seq( (4, "aa", "01/01/2017"), (2, "bb", "03/01/2017"), (6, "cc", "04/01/2017"), (1, "bb", "05/01/2017"), (5, "bb", "09/01/2017"), (3, "aa", "02/01/2017") )).toDF("id", "account", "time") // 转换时间列为日期类型 val df = rawDF.withColumn("date", to_date(col("time"), "dd/MM/yyyy"))
2. 核心逻辑:自连接筛选有序ID对
通过自连接将同一账户的记录两两匹配,再通过日期条件筛选出"早时间ID"和"晚时间ID"的组合:
val resultDF = df.alias("left") .join(df.alias("right"), // 连接条件:同一账户,且左表日期早于右表日期 col("left.account") === col("right.account") && col("left.date") < col("right.date"), "inner") // 重命名列并选择需要的字段 .select( col("left.id").alias("id1"), col("right.id").alias("id2"), col("left.account") ) // 查看结果 resultDF.show()
运行后输出的结果就和你需要的完全一致:
+---+---+-------+ |id1|id2|account| +---+---+-------+ | 4| 3| aa| | 2| 1| bb| | 2| 5| bb| | 1| 5| bb| +---+---+-------+
补充说明
- 为什么
cc没有出现在结果里?因为它只有一条记录,没有同账户的其他记录可以配对,符合需求逻辑。 - 如果你的时间格式不是
dd/MM/yyyy,记得调整to_date函数里的格式参数。
内容的提问来源于stack exchange,提问作者walker
相关产品推荐
相关产品推荐

