JavaPairRDD中groupBy分组逻辑、运行样例及返回值类型疑问
问题解答
1. groupBy分组逻辑样例演示
groupBy(row -> row.getString(row.fieldIndex("key")))的核心逻辑是:为输入的每一个Row对象生成分组键,Spark会将所有键相同的Row归类到同一个分组中。
以下是直观的运行样例:
测试数据(datasetOne.txt)
key,name,age 1,张三,22 2,李四,25 1,王五,23 3,赵六,24 2,孙七,21
执行过程
- 读取csv时因为配置了
option("header", "true"),第一行会被识别为字段名,后续每一行都会被封装为Row对象 - 对每个Row执行lambda表达式:
- 先调用
row.fieldIndex("key")获取key字段在Row中的列索引,这里返回索引值0 - 再调用
row.getString(0)拿到当前行key字段的实际值,作为当前行的分组键
- 先调用
- 所有行生成的分组键依次为:
"1"、"2"、"1"、"3"、"2"
分组结果
最终生成的JavaPairRDD包含3个键值对:
- 键
"1":对应2个Row对象:["1", "张三", 22]、["1", "王五", 23] - 键
"2":对应2个Row对象:["2", "李四", 25]、["2", "孙七", 21] - 键
"3":对应1个Row对象:["3", "赵六", 24]
2. 返回值中Iterable是否为List集合
答案是否定的,Iterable<Row>只是一个可迭代接口,不是具体的List实现类。
Spark不会将同一个key关联的所有条目一次性加载到内存的List中存储,该设计是为了避免单个key对应数据量过大时,全量加载到内存引发OOM异常。你可以通过Iterable的迭代器遍历该key下的所有Row,若确实需要转为List存储,可以手动调用工具类转换,但需要先确认单个key对应的数据量不会超出内存上限。
内容的提问来源于stack exchange,提问作者FreshTransistor
相关产品推荐
相关产品推荐

