【问题标题】:How to improve performance for slow Spark jobs using DataFrame and JDBC connection?如何使用 DataFrame 和 JDBC 连接提高慢速 Spark 作业的性能?
【发布时间】:2015-11-18 05:31:52
【问题描述】:

我正在尝试在单个节点(本地 [*])上以独立模式通过 JDBC 访问中型 Teradata 表(约 1 亿行)。

我使用的是 Spark 1.4.1。并且设置在非常强大的机器上(2 cpu、24 核、126G RAM)。

我尝试了几种内存设置和调整选项以使其工作更快,但它们都没有产生巨大的影响。

我确信我缺少一些东西,下面是我的最后一次尝试,它花了大约 11 分钟来获得这个简单的计数,而使用通过 R 的 JDBC 连接只用了 40 秒来获得计数。

bin/pyspark --driver-memory 40g --executor-memory 40g

df = sqlContext.read.jdbc("jdbc:teradata://......)
df.count()

当我尝试使用 BIG 表(5B 条记录)时,查询完成后没有返回任何结果。

【问题讨论】:

  • 使用 R 怎么算?
  • @zero323 - 在使用 Teradata JARS 设置连接后,只需使用 RJDBCteradataR 包..然后tdQuery("SELECT COUNT(*) FROM your_table)
  • 据我所知 Spark JDBC 数据源可以下推谓词,但实际执行是在 Spark 中完成的。这意味着您必须将数据传输到 Spark 集群。所以它与通过 JDBC(R 案例)执行 SQL 查询不同。首先你应该做的是在加载后缓存你的数据。但它不会提高第一次查询的性能。
  • @zero323 - 谢谢,在对此进行了更多研究后,我意识到这一点。我确实有一个快速的问题想 - 在 apache spark 中读取数据的最快方法是什么?是通过 Parquet 文件结构吗?
  • 这可能是一个不错的选择,但在使用这种方式之前,您可以尝试的第一件事是使用Teradata Hadoop conector。看起来它支持多个导出选项,包括 Hive 表。但是,使用单机网络和磁盘 IO 仍然是一个限制因素。

标签: apache-spark teradata pyspark spark-dataframe


【解决方案1】:

所有的聚合操作都是在整个数据集被检索到内存到DataFrame 集合中之后执行的。因此,在 Spark 中进行计数永远不会像直接在 TeraData 中那样高效。有时通过创建视图然后使用 JDBC API 映射这些视图来将一些计算推送到数据库中是值得的。

每次使用 JDBC 驱动程序访问大表时,都应指定分区策略,否则您将创建一个带有单个分区的 DataFrame/RDD,并且会使单个 JDBC 连接过载。

相反,您想尝试以下 AI(自 Spark 1.4.0+ 起):

sqlctx.read.jdbc(
  url = "<URL>",
  table = "<TABLE>",
  columnName = "<INTEGRAL_COLUMN_TO_PARTITION>", 
  lowerBound = minValue,
  upperBound = maxValue,
  numPartitions = 20,
  connectionProperties = new java.util.Properties()
)

还有一个选项可以下推一些过滤。

如果您没有均匀分布的整数列,您希望通过指定自定义谓词(where 语句)来创建一些自定义分区。例如,假设您有一个时间戳列并希望按日期范围进行分区:

    val predicates = 
  Array(
    "2015-06-20" -> "2015-06-30",
    "2015-07-01" -> "2015-07-10",
    "2015-07-11" -> "2015-07-20",
    "2015-07-21" -> "2015-07-31"
  )
  .map {
    case (start, end) => 
      s"cast(DAT_TME as date) >= date '$start'  AND cast(DAT_TME as date) <= date '$end'"
  }

 predicates.foreach(println) 

// Below is the result of how predicates were formed 
//cast(DAT_TME as date) >= date '2015-06-20'  AND cast(DAT_TME as date) <= date '2015-06-30'
//cast(DAT_TME as date) >= date '2015-07-01'  AND cast(DAT_TME as date) <= date '2015-07-10'
//cast(DAT_TME as date) >= date '2015-07-11'  AND cast(DAT_TME as date) <= date //'2015-07-20'
//cast(DAT_TME as date) >= date '2015-07-21'  AND cast(DAT_TME as date) <= date '2015-07-31'


sqlctx.read.jdbc(
  url = "<URL>",
  table = "<TABLE>",
  predicates = predicates,
  connectionProperties = new java.util.Properties()
)

它将生成一个DataFrame,其中每个分区将包含与不同谓词关联的每个子查询的记录。

DataFrameReader.scala查看源代码

【讨论】:

  • @zero323, @Gianmario Spacagna 如果我真的需要阅读整个MySQL 表(而不仅仅是获取count),那么如何提高Spark-SQLsluggish 性能呢?我已经使用spark.read.jdbc(..numPartitions..) 方法并行化读取操作。
  • 我的 MySQL (InnoDB) 表有 ~ 186M 记录,重约 149 GB(根据 phpMyAdmin 显示的统计数据)我正在使用numPartitions = 32。 [Spark 2.2.0] 我在 EMR 5.12.0 上,有 1 个 master、1 个 task 和 1 个 core(所有 r3.xlarge、8 个 @987654344 @,30.5 GiB 内存,80 SSD GB 存储)。我发现如果我不将 limit 记录到 ~ 1.5-2M,则将 MySQL 表读入 DataFrame 会失败。它提供了一个很长的 stack-trace,其中包含 javax.servlet.ServletException: java.util.NoSuchElementException: None.getjava.sql.SQLException: Incorrect key file for table..
【解决方案2】:

未序列化的表是否适合 40 GB?如果它开始在磁盘上进行交换,性能将急剧下降。

无论如何,当您使用带有 ansi SQL 语法的标准 JDBC 时,您会利用 DB 引擎,因此,如果 teradata(我不知道 teradata)保存有关您的表的统计信息,那么经典的“从表中选择计数(*)”将是非常快。 相反,spark 正在使用“select * from table”之类的内容将 1 亿行加载到内存中,然后将对 RDD 行进行计数。这是一个完全不同的工作量。

【讨论】:

  • 我认为它会,我也尝试将内存增加到 100 GB,但没有看到任何改进。我不是要在内存中加载 1 亿行,而是在数据帧上运行一些聚合操作,例如 count() 或临时表上的 count(*),但 Spark 花费的时间太长。我还尝试将 DF 注册为临时表并进行了简单的计数,但需要大约相同的时间。 ra1.registerTempTable("ra_dt"); total = sqlContext.sql("select count(*) from ra_dt")
  • 是的,但我认为 spark 不会修剪 DB 引擎上的计数操作,因此它会将所有行加载到内存中,然后对 DF 执行计数。
  • 该表中有多少列,1 亿行很容易达到 100 GB 的未序列化对象。你能发布你的表格架构吗?
  • 我认为你是对的,我在网上阅读了其他几篇文章,发现 Spark 正在尝试在应用计数操作之前加载数据。在这种情况下,在 Spark 中更快地读取此类数据的理想方法是什么?换句话说在 apache spark 中读取数据的最快方法是什么?这是我的表架构:root |-- field1: decimal(18,0) (nullable = true) |-- field2 : string (nullable = true) |-- field3: date (nullable = true) |-- field4: date (nullable = true) |-- field5: integer (nullable = true) |-- field6: string (nullable = true )
  • Spark 是一个分布式处理引擎,因此在 spark 中加载数据的最佳方式是从分布式文件系统或 dbms 中加载数据。在您的情况下,在一个单一实例上工作,我认为您只能通过指定 partitionColumn、lowerBound、upperBound、numPartition 来提高性能以提高读取并行性。如果您需要在计数之后执行其他查询,您可以在计数之前缓存 DF,因此第一次计数将花费时间,但接下来的查询将在内存中并且会更快。
【解决方案3】:

与其他解决方案不同的一个解决方案是将 oracle 表中的数据保存在 avro 文件中(在许多文件中分区)保存在 hadoop 上。 这种方式用 spark 读取这些 avro 文件会很轻松,因为你不会再调用 db。

【讨论】:

    猜你喜欢
    • 2011-09-26
    • 2017-06-25
    • 1970-01-01
    • 1970-01-01
    • 2017-02-06
    • 1970-01-01
    • 1970-01-01
    • 2020-09-09
    相关资源
    最近更新 更多