Spark词频统计代码对应DAG图的三个技术疑问
Spark词频统计DAG相关疑问解答
代码示例
spark = SparkSession.builder.appName("test").getOrCreate() sc = spark.sparkContext rdd = sc.textFile("data/words.txt") rdd1 = rdd.flatMap(lambda x: x.split(", ")) rdd2 = rdd1.map(lambda x: (x, 1)) rdd2.cache() rdd3 = rdd2.reduceByKey(lambda x, y: x + y) print(rdd3.collect()) rdd2.unpersist() time.sleep(1111111)
技术疑问
- 顶部蓝色块中textFile包含两个步骤,原因是什么?这两个步骤分别是什么?
- 代码中rdd2缓存前需经过flatMap和map处理,为何在DAG图中看不到这两个操作?
- reduceByKey包含两个步骤,这两个步骤分别是什么?为何它们不在同一个蓝色块中?
问题解答
问题1
textFile对应的两个步骤分别是HadoopRDD和MapPartitionsRDD,这是Spark读取文件的底层逻辑决定的:
HadoopRDD:负责对接Hadoop InputFormat,从文件系统读取原始字节流,是数据读取的底层操作。MapPartitionsRDD:将读取到的字节流按行解析成文本行,完成原始数据到可处理文本的转换。
这两步是Spark文件读取的默认流程,因此在DAG中会拆分为两个步骤显示。
问题2
flatMap和map都属于窄依赖操作,且和textFile的读取操作同属一个Stage。Spark的DAG可视化会把连续的窄依赖操作合并到同一个蓝色块(同一Task集合)中展示,不会单独列出每个窄依赖转换。你看到的包含textFile的蓝色块,实际上已经整合了后续的flatMap和map操作。
问题3
reduceByKey的两个步骤分别是Map端本地聚合和Reduce端全局聚合:
- Map端本地聚合:在每个分区内对相同key的value做预聚合,减少Shuffle阶段的数据传输量,属于Shuffle前的Stage 0。
- Reduce端全局聚合:经过Shuffle将相同key的数据拉到同一节点后,完成最终的聚合,属于Shuffle后的Stage 1。
它们不在同一个蓝色块的核心原因是Shuffle操作是Stage划分的边界,Spark会以Shuffle为分割点拆分Stage,因此两个步骤分属不同的Stage,对应不同的蓝色块。
内容的提问来源于stack exchange,提问作者user3692015
相关产品推荐
相关产品推荐

