Spark读写MySQL时的任务数量及写入连接数技术咨询
Spark读写MySQL的任务数与连接数解析
一、读取MySQL时的任务数与数据分配逻辑
Spark读取MySQL的任务数由**分区数(numPartitions)**决定,具体分两种场景:
- 若指定了
numPartitions参数(如第一个示例代码):Spark会创建与分区数相等的任务,每个任务仅读取数据的一个分片(而非全量)。它会通过拆分查询条件(需配合partitionColumn、lowerBound、upperBound参数)将数据划分为多个区间,每个任务负责对应区间的数据读取,实现并行读取。 - 若未指定
numPartitions(如第二个示例代码):默认仅创建1个任务,由单个Executor读取全量数据后,再将数据分发至集群其他节点处理,这种方式无法利用集群并行能力,仅适合小数据量场景。
二、写入MySQL时的连接数
Spark写入MySQL的连接数与当前作业的并行任务数完全一致:
- 每个写入任务会单独建立一个JDBC连接到MySQL,即并行任务数等于连接数。
- 并行任务数通常由DataFrame的分区数决定,而分区数取决于前期处理逻辑(如读取时的分区数、shuffle后的分区数等)。
- 若需控制连接数避免压垮MySQL,可通过
repartition或coalesce调整DataFrame的分区数,减少并行写入任务量。
示例代码说明
带分区配置的读取函数
def jdbc(sqlContext: SQLContext, url: String, driver: String, dbtable: String, user: String, password: String, numPartitions: Int): DataFrame = { sqlContext.read.format("jdbc").options(Map( "url" -> url, "driver" -> driver, "dbtable" -> s"(SELECT * FROM $dbtable) $dbtable", "user" -> user, "password" -> password, "numPartitions" -> numPartitions.toString )).load }
该函数通过numPartitions指定读取并行度,配合分区列参数可实现多任务并行分片读取,提升大数据量读取效率。
默认配置的读取函数
def mysqlToDF(sparkSession:SparkSession, jdbc:JdbcInfo, table:String): DataFrame ={ var dF1 = sparkSession.sqlContext.read.format("jdbc") .option("url", jdbc.jdbcUrl) .option("user", jdbc.user) .option("password", jdbc.passwd) .option("driver", jdbc.jdbcDriver) .option("dbtable", table) .load() dF1.createOrReplaceTempView(s"${table}") dF1 }
该函数未指定分区参数,默认单任务读取全量数据,仅适合小数据集场景,并行处理能力不足。
内容的提问来源于stack exchange,提问作者Shengxin Huang
相关产品推荐
相关产品推荐

