PySpark中对数组的数组列进行双层洗牌并提取指定元素
处理PySpark嵌套数组的双层洗牌与元素提取
问题背景
我有一个PySpark列gm_array,格式如下:
gm_array [[1, 4, 6,...], [2, 7, 8,...], [3, 5, 7,...],...] [[8, 11, 9,...], [7, 2, 6,...], [10, 9, 8,...],...] [[90, 13, 67,...], [55, 6, 98,...], [1, 6, 2,...],...] . .
需要完成两个操作:
- 对
gm_array的外层数组和每个内层数组分别进行随机洗牌,得到洗牌后的完整数组列; - 从洗牌后的数组中提取前5个内层数组的第一个元素,组成新数组。
解决方案
1. 导入依赖模块
from pyspark.sql import SparkSession from pyspark.sql.functions import udf, col from pyspark.sql.types import ArrayType, IntegerType import random
2. 定义双层洗牌UDF
通过自定义UDF实现外层数组和内层数组的双重随机洗牌:
def shuffle_nested_array(arr): # 洗牌外层数组 shuffled_outer = random.sample(arr, len(arr)) # 逐个洗牌内层数组 shuffled_inner = [random.sample(inner_arr, len(inner_arr)) for inner_arr in shuffled_outer] return shuffled_inner # 注册UDF,若内层元素为其他类型(如浮点型),需替换IntegerType为对应类型 shuffle_nested_udf = udf(shuffle_nested_array, ArrayType(ArrayType(IntegerType())))
3. 生成洗牌后的数组列
将UDF应用到原DataFrame,得到洗牌后的列:
# 假设原DataFrame名为df df_shuffled = df.withColumn("shuffled_gm_array", shuffle_nested_udf(col("gm_array")))
4. 提取前5个内层数组的首元素
定义UDF提取目标元素:
def extract_first_elements(arr): # 取前5个内层数组,每个数组取第一个元素 return [inner_arr[0] for inner_arr in arr[:5]] extract_first_udf = udf(extract_first_elements, ArrayType(IntegerType())) # 添加提取结果列 df_result = df_shuffled.withColumn("first_elements", extract_first_udf(col("shuffled_gm_array")))
查看输出
查看洗牌后的完整数组列:
df_result.select("shuffled_gm_array").show(truncate=False)示例输出:
shuffled_gm_array [[19, 6, 1,...], [9, 80, 5,...], [30, 7, 3,...],...] [[7, 9, 11,...], [6, 8, 7,...], [18, 7, 10,...],...] [[90, 1, 7,...], [8, 9, 81,...], [6, 5, 1,...],...] . .查看提取的首元素数组列:
df_result.select("first_elements").show(truncate=False)示例输出:
[19, 9, 30,...] [7, 6, 18,...] [90, 8, 6,...] . .
注意事项
- 若数组元素类型非整数,需调整UDF中的数据类型(如
DoubleType()对应浮点型); - 如需固定随机种子保证结果可复现,可在
shuffle_nested_array函数内添加random.seed(固定数值),但会导致所有分区洗牌结果一致,按需选择。
内容的提问来源于stack exchange,提问作者Chris_007
相关产品推荐
相关产品推荐

