【问题标题】:How many tasks are created when spark read or write from mysql?spark从mysql读取或写入时创建了多少个任务?
【发布时间】:2022-10-07 01:35:20
【问题描述】:

据我所知,Spark executors同时处理多个任务,保证并行处理数据。问题来了。当连接到外部数据存储时,比如说mysql,有多少任务可以完成这个工作?换句话说,是同时创建多个任务,每个任务读取所有数据,还是只从一个任务中读取数据并分发以其他方式到集群?向mysql写入数据怎么样,有多少个连接?

这是一些从/向mysql读取或写入数据的代码:


    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
  }

  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.show(3)
    dF1.createOrReplaceTempView(s\"${table}\")
    dF1

  }
}

    标签: mysql apache-spark


    【解决方案1】:

    这是一篇很好的文章,可以回答您的问题: https://freecontent.manning.com/what-happens-behind-the-scenes-with-spark/

    简单来说:worker 将读取任务分成几个部分,每个 worker 只读取你输入数据的一部分。划分的任务数量取决于您的资源和数据量。写入原理相同:Spark 将数据写入分布式存储系统,例如 Hdfs,而在 Hdfs 中数据以分布式方式存储:每个 worker 将其数据写入 Hdfs 中的某个存储节点。

    【讨论】:

      【解决方案2】:

      默认情况下,来自 jdbc 源的数据由一个线程加载,因此您将有一个任务由一个执行程序处理,这就是您在第二个函数 mysqlToDF 中可能期望的情况

      在第一个函数“jdbc”中,您更接近并行读取,但仍然需要一些参数,numPartitions 是不够的,spark 需要一些整数/日期列和下/上限才能并行读取(它将执行 x 个查询部分结果)

      Spark jdb documentation

      在本文档中,您会发现:

      partitionColumn、lowerBound、upperBound(无)这些选项必须 如果指定了其中任何一个,则全部指定。此外, 必须指定 numPartitions。他们描述了如何划分 从多个工作人员并行读取时的表。分区列 必须是表中的数字、日期或时间戳列 问题。请注意,lowerBound 和 upperBound 仅用于 决定分区步长,而不是用于过滤表中的行。所以 表中的所有行都将被分区并返回。这个选项 仅适用于阅读。

      numPartitions(无)最大值 可用于表读取的并行性的分区数 和写作。这也决定了最大并发数 JDBC 连接。如果要写入的分区数超过这个 限制,我们通过调用 coalesce(numPartitions) 将其减少到这个限制 在写作之前。读/写

      关于写

      向mysql写入数据怎么样,有多少个连接?

      如文档中所述,它还取决于 numPartitions,如果写入时的分区数高于 numPartitions,Spark 会计算出来并调用 coalesce。请记住,合并可能会产生偏差,因此有时最好使用 repartition(numPartitions) 显式重新分区以在写入之前平均分配数据

      如果您未设置 numPartitions 写入时的并行连接数可能与给定时刻的活动任务数相同,因此请注意,如果并行度高且没有上限,您可能会阻塞源服务器

      【讨论】:

        猜你喜欢
        • 2022-12-01
        • 2020-07-14
        • 2017-03-28
        • 1970-01-01
        • 1970-01-01
        • 2021-08-22
        • 2018-12-31
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多