Spark集群中写入数据库的操作由Driver还是Executor执行?
Spark写入MySQL与文本文件的执行节点说明
写入MySQL的执行节点
写入MySQL的操作分两种场景,执行主体不同:
- 推荐的分布式写入方式(Spark JDBC API):由Executor执行
当使用Spark原生的DataFrameWriter.jdbc()(或RDD对应的JDBC写入方式)时,数据分布在各个Executor的分区中,每个Executor会独立建立JDBC连接,将自己负责的分区数据写入MySQL。Driver仅负责下发写入逻辑,实际写入操作由各个Executor分布式完成,这种方式效率高,适合大数据量场景。 - Driver端收集后写入:由Driver执行
如果先调用collect()/toPandas()将分布式数据拉取到Driver端,再用Python的pymysql等客户端库写入,此时写入操作由Driver执行。但这种方式仅适合极小数据量,大数据量会导致Driver内存溢出,生产环境不推荐。
关于saveAsTextFile的误解纠正
你之前认为saveAsTextFile由Driver执行是错误的,这个操作由各个Executor分布式完成:
- 执行时,每个Executor会将自己负责的RDD分区数据写入文件,最终输出的是一个目录(而非单个文件),目录下包含多个
part-xxxxx文件(每个分区对应一个),还有_SUCCESS标记文件。 - 若你指定的是本地路径,每个Worker节点的Executor会在自身节点的本地文件系统生成对应文件;如果是HDFS等分布式存储路径,才会统一写入到分布式系统中。你可能在Driver所在节点看到了输出目录,但这只是因为Driver节点可能同时也是Worker节点,并非Driver执行了写入操作。
你的文本文件写入代码:
from pyspark import SparkConf,SparkContext if __name__ == "__main__": sc = SparkContext(master = "spark://IP:PORT",appName='word_count_application') words = sc.textFile("book_2.txt") word_count = words.flatMap(lambda a : a.split(" ")).map(lambda a : (a,1)).reduceByKey(lambda a,b : a+b) word_count.saveAsTextFile("book2_output.txt")
内容的提问来源于stack exchange,提问作者Saranraj K
相关产品推荐
相关产品推荐

