PySpark按多条件创建新列,统计Y值在对应X时间前的出现次数
PySpark 实现方案
核心思路
- 基于Spark窗口函数完成计算,不会修改原数据集的物理行顺序,仅在窗口计算逻辑内按X时间排序统计,完全符合不排序原X列的要求。
实现代码
1. 导入依赖
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import col, count, lit
2. 初始化Spark与构造测试数据(实际使用可替换为读取自己的数据集)
# 初始化SparkSession spark = SparkSession.builder.appName("count_history_occur").getOrCreate() # 示例输入数据 data = [ ("2021-09-08", "number1"), ("2021-09-09", "number2"), ("2021-09-10", "number2"), ("2021-09-11", "number3"), ("2021-09-12", "number2"), ("2021-09-13", "number2"), ("2021-09-14", "number3") ] df = spark.createDataFrame(data, schema=["X", "Y"]) # 转换X列为timestamp类型 df = df.withColumn("X", col("X").cast("timestamp"))
3. 核心计算逻辑
# 定义窗口规则:按Y值分区,分区内按X时间升序排列,统计范围为分区开头到当前行的前一行 window_spec = Window.partitionBy("Y").orderBy("X").rowsBetween(Window.unboundedPreceding, -1) # 新增Z列 df_result = df.withColumn("Z", count(lit(1)).over(window_spec).cast("int"))
4. 结果验证
执行df_result.show()即可得到期望的输出结果,原数据集的行顺序不会发生任何改变。
逻辑说明
窗口按Y字段分组,保证仅统计和当前行相同Y值的记录;窗口内的orderBy仅作用于分区内的计算逻辑,不会对整个数据集做全局排序;帧范围设置为从分区第一条到当前行的前一条,刚好统计所有X时间早于当前行的同Y值记录数,空窗口的count结果自动为0,完全匹配需求。
内容的提问来源于stack exchange,提问作者Kerem Tatlıcı
相关产品推荐
相关产品推荐

