fetchsize 和 batchsize 对 Spark 的影响

Sco*_*ieh 7 database performance apache-spark spark-dataframe

我想直接通过Spark控制RDB的读写速度,但是标题已经显示的相关参数似乎不起作用。

我可以得出结论,fetchsize并且batchsize不能使用我的测试方法吗?或者它们确实影响阅读和写作的方面,因为基于规模的衡量结果是合理的。

的统计betchsizefetchsize数据和设置

/*Dataset*/
+--------------+-----------+
| Observations | Dataframe |
+--------------+-----------+
|      109,077 | Initial   |
|      345,732 | Ultimate  |
+--------------+-----------+
/*fetchsize*/
+-----------+-----------+------------------+------------------+
| fetchsize | batchsize | Reading Time(ms) | Writing Time(ms) |
+-----------+-----------+------------------+------------------+
|        10 |        10 |            2,103 |           38,428 |
|       100 |        10 |            2,123 |           38,021 |
|     1,000 |        10 |            2,032 |           38,345 |
|    10,000 |        10 |            2,016 |           37,892 |
|    50,000 |        10 |            2,017 |           37,795 |
|   100,000 |        10 |            2,055 |           38,720 |
+-----------+-----------+------------------+------------------+
/*batchsize*/
+-----------+-----------+------------------+------------------+
| fetchsize | batchsize | Reading Time(ms) | Writing Time(ms) |
+-----------+-----------+------------------+------------------+
|        10 |        10 |            2,072 |           37,977 |
|        10 |       100 |            2,077 |           36,990 |
|        10 |     1,000 |            2,034 |           36,703 |
|        10 |    10,000 |            1,979 |           36,980 |
|        10 |    50,000 |            2,043 |           36,749 |
|        10 |   100,000 |            2,005 |           36,624 |
+-----------+-----------+------------------+------------------+
Run Code Online (Sandbox Code Playgroud)

Datadog观察的措施

MySQL 性能测量

可能有帮助的细节

我在AWS上创建了两个m4.xlarge Linux 实体,一个用于Spark的执行,另一个用于 RDB 上的数据存储,使用Datadog来观察Spark应用程序的性能,尤其是对 RDB 的读写。Spark处于独立模式,测试应用程序只是从MySQL RDB 中提取一些数据,进行一些计算,然后推回MySQL

一些细节如下:

  1. JDBC 属性放在文件application.conf 中,如下所示:

    spark {
      Reading {
        url: "jdbc:mysql://address/designated database"
        driver: "com.mysql.cj.jdbc.Driver"
        user: "username"
        password: "password"
        fetchsize: "10000"
      }
      Writing {
        url: "jdbc:mysql://address/designated database"
        driver: "com.mysql.cj.jdbc.Driver"
        dbtable: "designated table"
        user: "username"
        password: "password"
        batchsize: "10000"
        truncate: "true"
      }
    }
    
    Run Code Online (Sandbox Code Playgroud)
  2. 执行应用程序时记录日志由log4jx2启用,在其中,测量写入时间。

                    .
                    .
                    .
    startTime = System.nanoTime()
    val connection = new Properties()
    configureProperties(connection, conf, "spark.Writing")
    val ultimateObservations = ultimateResult.count()
    ultimateResult.write
        .mode(SaveMode.Overwrite)
        .jdbc(conf.getString("spark.Writing.url"),
              conf.getString("spark.Writing.dbtable"),
              connection)
    finishedTime = System.nanoTime()
    logger.info("Finished writing from Spark to MySQL, taking {} milliseconds; approximately {} rows/s",
          TimeUnit.MILLISECONDS.convert((finishedTime - startTime), TimeUnit.NANOSECONDS),
          ultimateObservations/TimeUnit.SECONDS.convert((finishedTime - startTime), TimeUnit.NANOSECONDS)
        )
                    .
                    .
                    .
    
    /*
     *configureProperties is a customized function
     */
    def configureProperties(connectionEntity: Properties, conf: Config, designatedString: String): Unit = {
        val propertiesCarrier = conf.getConfig(designatedString)
        for (entry <- propertiesCarrier.entrySet) {
          if (entry.getKey().trim() != "url" && entry.getKey().trim() != "dbtable") {
            connectionEntity.put(entry.getKey(), entry.getValue().unwrapped().toString())
            logger.info("Database configuration: ({}, {}).",
      entry.getKey(), entry.getValue().unwrapped().toString: Any)
        }
      }
    }
    
    Run Code Online (Sandbox Code Playgroud)