You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将含列表的元组扁平化为多个元组(RDD场景)

Spark RDD 键值对扁平化处理方案

要实现将每个(key, List[value])结构的RDD元素拆分为多个(key, single_value)的元组,核心是使用Spark的flatMap算子,它可以完成先映射再扁平化的操作。

Python 实现示例

假设你已经创建了原始RDD:

from pyspark import SparkContext
sc = SparkContext("local", "FlattenExample")
original_rdd = sc.parallelize([('ABE', ['ORD', 'ATL', 'DTW'])])

执行扁平化处理:

flattened_rdd = original_rdd.flatMap(lambda item: [(item[0], elem) for elem in item[1]])
# 输出结果
print(flattened_rdd.collect())

执行后会得到期望输出:

[('ABE', 'ORD'), ('ABE', 'ATL'), ('ABE', 'DTW')]

Scala 实现示例

原始RDD定义:

import org.apache.spark.SparkContext
import org.apache.spark.SparkConf

val conf = new SparkConf().setAppName("FlattenExample").setMaster("local")
val sc = new SparkContext(conf)
val originalRDD = sc.parallelize(Seq(("ABE", List("ORD", "ATL", "DTW"))))

扁平化处理代码:

val flattenedRDD = originalRDD.flatMap { case (key, values) => values.map(elem => (key, elem)) }
// 查看结果
flattenedRDD.collect().foreach(println)

输出结果同样符合要求:

(ABE,ORD)
(ABE,ATL)
(ABE,DTW)

原理说明

flatMap会对RDD中的每个元素执行两步操作:

  1. 映射:将单个(key, 列表)元组转换为一个包含多个(key, 元素)元组的列表
  2. 扁平化:将所有生成的列表展开,把列表内的元组作为独立元素存入新的RDD

内容的提问来源于stack exchange,提问作者Solomon123

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.01 14:40:28