PySpark RDD内部数组操作:如何在Worker节点中对元素乘2?
解决Spark RDD中子数组元素批量乘2的问题
你可以通过Spark的map算子直接在Worker节点上处理每个RDD元素,无需将数据拉回Driver端(比如take这类会把数据取回Driver的操作)。map会把处理逻辑分发到各个Worker节点,针对每个分区内的元组独立执行操作。
具体实现代码
sc = SparkContext() x = [(1, [2, 3, 4, 5]), (2, [2, 7, 8, 10])] y = sc.parallelize(x) # 用map处理每个元组:保留第一个元素,对第二个数组的每个元素乘2 processed_rdd = y.map(lambda t: (t[0], [num * 2 for num in t[1]])) # 验证处理结果 print(processed_rdd.collect())
代码说明
lambda t: (t[0], [num * 2 for num in t[1]]):这个匿名函数针对RDD中的每个元组t,保留元组的第一个元素t[0],同时遍历第二个元素(数组)的每个元素,执行乘2操作生成新数组,最终组成新的元组。map算子会将这段处理逻辑推送到各个Worker节点,每个节点仅处理自己负责的RDD分区数据,全程不需要将分布式数据拉回Driver,完全符合Spark的分布式处理逻辑。
执行代码后会输出预期结果:
[(1, [4, 6, 8, 10]), (2, [4, 14, 16, 20])]
内容的提问来源于stack exchange,提问作者AutoOffice
相关产品推荐
相关产品推荐

