Spark为每行生成1-30范围内随机数的最优实现方案
解决Spark DataFrame每行生成1-30随机数并避免数据倾斜的问题
原代码的问题
你用的lit(Random.nextInt(30))会在Driver端生成一个固定随机数,然后将这个值广播到所有Executor,最终DataFrame的所有行都会使用同一个随机数,不符合“每行生成不同随机数”的需求。
正确实现方式
要在分布式环境下为每行生成独立的1-30范围内随机数,需使用Spark内置的分布式随机函数rand(),它会在Executor端为每行生成独立的随机值:
import org.apache.spark.sql.functions._ // 生成1-30的整数随机数 val dfWithRandom = df.withColumn("random_num", floor(rand() * 30) + 1)
代码说明
rand():生成0到1之间的均匀分布随机浮点数,每个Executor独立计算,适合70亿行的大数据场景。floor(rand() * 30):将0-1的随机数放大30倍后取整,得到0-29的整数。+1:将范围偏移为1-30,符合需求。
避免数据倾斜的关键
因为后续要基于random_num做repartition,需保证随机数分布均匀:
rand()生成的是均匀分布的随机值,1-30每个数字的出现概率基本一致,以此字段做分区时,各分区的数据量会相对均衡,从根源避免数据倾斜。- 如果需要固定随机结果(用于调试或重复计算),可以给
rand()传入种子参数,比如rand(1234L),这样每次运行生成的随机数序列一致。
内容的提问来源于stack exchange,提问作者HEMANT PATEL
相关产品推荐
相关产品推荐

