如何为无唯一值列的Spark DataFrame添加连续行号列?
没问题,我来帮你搞定这个生成连续行号的需求!你遇到的monotonically_increasing_id()的问题很常见——它生成的是跨分区的递增ID,每个分区的ID区间是独立的,所以肯定不会是连续的1、2、3。下面给你两种可行的解决方案,一种用窗口函数(解决你的分区顾虑),另一种完全不用窗口函数:
方法一:窗口函数 + 强制单分区
你担心窗口函数因为分区问题生成不连续的行号?那我们可以先把DataFrame合并到单个分区,再用row_number()生成行号。这样就能保证所有数据在同一个分区里,行号自然连续:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 先将数据合并到1个分区,再添加连续行号 age_with_id = age.repartition(1).withColumn( "col_id", F.row_number().over(Window.orderBy(F.monotonically_increasing_id())) ) age_with_id.show()
输出就是你想要的:
+---+------+ |age|col_id| +---+------+ | 10| 1| | 11| 2| | 13| 3| +---+------+
这里orderBy(F.monotonically_increasing_id())只是为了给窗口函数一个排序依据(Spark要求窗口函数必须有order by),因为已经是单分区了,行的顺序会和原DataFrame一致(如果原数据没经过shuffle的话)。
方法二:用RDD的zipWithIndex(无需窗口函数)
如果你不想用窗口函数,可以把DataFrame转成RDD,用zipWithIndex给每个元素分配连续索引,再转回DataFrame。注意zipWithIndex是从0开始的,所以要加1得到1起始的行号:
# 转成RDD添加索引,再转回DataFrame age_rdd = age.rdd.zipWithIndex() # 索引+1,转换成(age, col_id)的结构 age_with_id = age_rdd.map(lambda x: (x[0][0], x[1] + 1)).toDF(["age", "col_id"]) age_with_id.show()
这个方法更直接,不需要处理分区问题,生成的行号绝对连续。
注意事项
- 两种方法对于小数据量都很友好,但如果是超大数据集,
repartition(1)或者RDD转换可能会有性能瓶颈(因为所有数据都要挤在一个节点上)。不过你的示例数据很小,完全不用担心。 - Spark本身不保证DataFrame的默认顺序,如果你需要严格按照某个逻辑排序后生成行号,把
orderBy里的参数换成你需要排序的列即可(比如Window.orderBy("age"))。
内容的提问来源于stack exchange,提问作者Metadata
相关产品推荐
相关产品推荐

