如何在Spark中按原始顺序克隆列值
解决Spark数据集按原始顺序整体克隆n次的问题
我完全懂你的需求——你想要把数据集按原始顺序完整重复n次,而不是把每个元素单独复制多份后排列。用flatMap确实会把每个元素先重复n次,导致顺序变成1,1,2,2,3,3,这里给你两个实用的解决方案:
方法一:使用Union重复合并数据集
这个方法简单直接,适合中小规模的数据集:
%spark val ds = spark.range(1, 4).withColumnRenamed("id", "v") val cloneCount = 2 // 生成包含n份原数据集的序列,再依次合并 val clonedDs = (1 to cloneCount).map(_ => ds).reduce(_ union _) clonedDs.show()
运行后会得到你想要的结果:
+---+ | v| +---+ | 1| | 2| | 3| | 1| | 2| | 3| +---+
原理是:(1 to cloneCount).map(_ => ds)创建了一个包含cloneCount个原数据集的集合,reduce(_ union _)把这些数据集按顺序合并,完美保留原始元素的排列顺序。
方法二:使用笛卡尔积(适合大数据集)
如果你的数据集很大,多次union可能会产生多个执行阶段(Stage),影响效率。这时可以用笛卡尔积的方式,通过一个辅助数据集来实现整体重复:
%spark val ds = spark.range(1, 4).withColumnRenamed("id", "v") val cloneCount = 2 // 创建包含1到cloneCount的辅助数据集 val helperDs = spark.range(1, cloneCount + 1) // 笛卡尔积后只保留原列,可选按辅助列排序确保顺序 val clonedDs = ds.crossJoin(helperDs) .select("v") .orderBy(helperDs("id")) // 若需要严格保证顺序可添加,多数场景下不加也符合预期 clonedDs.show()
这个方法通过辅助数据集和原数据集做笛卡尔积,相当于把原数据集的每个元素和辅助集中的每个元素配对,最后只保留原列,就实现了整体重复n次的效果。而且笛卡尔积是单个执行阶段,大数据场景下效率更优。
内容的提问来源于stack exchange,提问作者yogesh.tewari
相关产品推荐
相关产品推荐

