如何在PySpark DataFrame中按分组生成随机列?
解决PySpark中分组生成随机数的窗口函数报错问题
错误原因
PySpark窗口函数仅支持确定性表达式,而rand()属于非确定性函数(每次调用返回不同结果),因此无法直接在over(Window.partitionBy())中使用,触发AnalysisException。
解决方案
根据需求,有两种可行方法:
方法1:分组关联随机种子生成独立随机序列
如果需要每个分组拥有独立的随机数序列(同一分组内每行随机数不同,不同分组序列独立),可以先为每个分组生成唯一随机种子,再基于种子生成随机数:
import pyspark.sql.functions as F # 1. 为每个分组生成唯一随机种子 group_seeds = df.select("groups").distinct().withColumn("seed", F.randint(0, 1000000)) # 2. 将种子关联回原DataFrame df_with_seed = df.join(group_seeds, on="groups", how="left") # 3. 基于种子生成随机数 result = df_with_seed.withColumn("random_groups", F.rand(F.col("seed"))).drop("seed") result.show()
方法2:直接生成全局随机数(无需窗口函数)
如果只是需要每行生成随机数,不强调分组间的序列独立性,直接调用rand()即可:
import pyspark.sql.functions as F result = df.withColumn("random_groups", F.rand()) result.show()
预期输出示例
运行上述代码后,将得到类似以下的结果:
ID | groups | random_groups 1 | A | 0.3 2 | A | 0.9 3 | B | 0.8
内容的提问来源于stack exchange,提问作者Sam Comber
相关产品推荐
相关产品推荐

