PySpark如何为满足条件的每行随机分配列表中的不同值
问题原因
你使用random.choice(my_list)会导致所有符合条件的行拿到同一个随机值,因为这是Python本地函数,只会在Driver端执行一次,生成的单个值会被复用给所有行。
解决方案
要用Spark的内置函数实现每行独立生成随机索引,再从列表中取值,这样每行都会得到不同的随机元素:
方法一:使用element_at(索引从1开始)
from pyspark.sql import functions as f my_list = ["值1", "值2", "值3"] # 替换成你的目标列表 condition = "指定条件值" # 替换成你的判断条件值 df = df.withColumn( "rand_col", f.when( f.col("condition_col") == condition, # 生成1到列表长度的随机整数作为索引,取对应元素 f.element_at(f.array(*[f.lit(item) for item in my_list]), f.floor(f.rand() * len(my_list)) + 1) ) )
方法二:使用getItem(索引从0开始)
from pyspark.sql import functions as f my_list = ["值1", "值2", "值3"] condition = "指定条件值" df = df.withColumn( "rand_col", f.when( f.col("condition_col") == condition, # 生成0到列表长度-1的随机整数作为索引,取对应元素 f.array(*[f.lit(item) for item in my_list]).getItem(f.floor(f.rand() * len(my_list))) ) )
核心说明
f.rand()会为每行生成一个0到1之间的随机浮点数,乘以列表长度后取整,得到合法的索引值f.array(*[f.lit(item) for item in my_list])把Python列表转换成Spark的数组列,供后续索引取值- 两种方法都会让每个符合条件的行独立生成随机元素,不会出现全局统一值的问题
内容的提问来源于stack exchange,提问作者Piotr
相关产品推荐
相关产品推荐

