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

如何将SAS的first.variable分组序号逻辑转换为PySpark实现

SAS逻辑对应PySpark实现方案

SAS中if first.loan then seq_id = 0; seq_id+1;的核心逻辑是:按loan字段分组后,识别每个分组的第一行,从0开始逐行累加1,最终得到组内从1开始的连续自增序号。


实现前提说明

PySpark为分布式计算引擎,数据不保留默认的全局物理顺序,必须显式指定组内排序规则才能得到稳定的序号结果。给出的样例中同一loan下module值递增,可直接用module作为组内排序字段;如果实际业务有其他排序依据(如交易时间、入库时间等),替换排序字段即可。


完整实现代码

from pyspark.sql import SparkSession
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number

# 初始化Spark会话
spark = SparkSession.builder.appName("seq_id_generate").getOrCreate()

# 构造样例数据集
source_data = [
    (743,455),
    (4490,795),
    (1101,235),
    (1101,335),
    (1101,435),
    (3471,898)
]
df = spark.createDataFrame(source_data, schema=["loan", "module"])

# 定义窗口规则:按loan分组,组内按指定字段排序
window_rule = Window.partitionBy("loan").orderBy("module")

# 生成自增seq_id字段
result_df = df.withColumn("seq_id", row_number().over(window_rule))

# 打印结果验证
result_df.show()

输出结果验证

代码运行后输出如下,完全匹配预期:

+----+------+------+
|loan|module|seq_id|
+----+------+------+
| 743|   455|     1|
|1101|   235|     1|
|1101|   335|     2|
|1101|   435|     3|
|3471|   898|     1|
|4490|   795|     1|
+----+------+------+

无明确排序字段的兼容方案

如果数据集没有可用于组内排序的业务字段,需要先添加全局单调递增ID标记原始数据顺序,再基于该ID生成组内序号,避免分布式计算导致的序号随机问题:

from pyspark.sql.functions import monotonically_increasing_id

# 标记原始数据顺序
df = df.withColumn("raw_order", monotonically_increasing_id())
window_rule = Window.partitionBy("loan").orderBy("raw_order")
result_df = df.withColumn("seq_id", row_number().over(window_rule)).drop("raw_order")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 15:36:29