【问题标题】:Spark SQL from raw text to Parquet: no performance boost从原始文本到 Parquet 的 Spark SQL:没有性能提升
【发布时间】:2019-01-05 07:36:40
【问题描述】:

场景如下:

我有一个SparkSQL 程序,它在几个Hive 表上执行ETL 过程。这些表是使用 Sqoop 在 RAW TEXT 中使用 Snappy 压缩从 Teradata 数据库导入的(不幸的是,Avro format 不适用于 Teradata 连接器)。 Spark SQL 进程完成所需的时间约为 1 小时 15 分钟

为了提高性能,我想在执行 SparkSQL 进程之前将表转换为更有效的格式,例如 Parquet。根据文档和在线讨论,这应该会极大地促进使用原始文本(即使使用 snappy 压缩,在原始文本上不可拆分)。 因此,我使用 Snappy 压缩将所有 Hive 表转换为 Parquet 格式。我已经使用相同的设置(num-executors、driver-memory、executor-memory)在这些表上启动了 SparkSQL 进程。该过程在 1 小时 20 分钟后结束。 这对我来说非常令人惊讶。我没想到会有 30 倍的提升,就像我在一些讨论中看到的那样,但我当然期待改进。

Spark程序中执行的操作类型大部分是join和filter(where条件),如下sn-p所示:

val sc = new SparkContext(conf)
val sqc = new HiveContext(sc)
sqc.sql("SET hive.exec.compress.output=true")
sqc.sql("SET parquet.compression=SNAPPY")


var vcliff = sqc.read.table(s"$swamp_db.DBU_vcliff")
var vtktdoc = sqc.read.table(s"$swamp_db.DBU_vtktdoc")
var vasccrmtkt = sqc.read.table(s"$swamp_db.DBU_vasccrmtkt")

val numPartitions = 7 * 16

// caching
vcliff.registerTempTable("vcliff")
vtktdoc.registerTempTable("vtktdoc")
vasccrmtkt.registerTempTable("vasccrmtkt")


ar ORI_TktVCRAgency = sqc.sql(
    s"""
       |            SELECT tic.CODCLI,
       |            tic.CODARLPFX,
       |            tic.CODTKTNUM,
       |            tic.DATDOCISS,
       |            vloc.CODTHR,
       |            vloc.NAMCMPNAMTHR,
       |            vloc.CODAGNCTY,
       |            vloc.NAMCIT,
       |            vloc.NAMCOU,
       |            vloc.CODCOU,
       |            vloc.CODTYPTHR,
       |            vloc.CODZIP,
       |            vcom.CODCOMORGLEVDPC,
       |            vcom.DESCOMORGLEVDPC,
       |            vcom.CODCOMORGLEVRMX,
       |            vcom.DESCOMORGLEVRMX,
       |            vcom.CODCOMORGLEVSALUNT,
       |            vcom.CODPSECOMORGCTYLEVSALUNT,
       |            vcom.DESCOMORGLEVSALUNT,
       |            vcom.CODCOMORGLEVRPR,
       |            vcom.CODPSECOMORGCTYLEVRPR,
       |            vcom.DESCOMORGLEVRPR,
       |            vcom.CODCOMORGLEVCTYCNL,
       |            vcom.CODPSECOMORGCTYLEVCTYCNL,
       |            vcom.DESCOMORGLEVCTYCNL,
       |            vcom.CODCOMORGLEVUNT,
       |            vcom.CODPSECOMORGCTYLEVUNT,
       |            vcom.DESCOMORGLEVUNT,
       |            vcli.DESCNL
       |            FROM $swamp_db.DBU_vlocpos vloc
       |                LEFT JOIN $swamp_db.DBU_vcomorghiemktgeo vcom ON vloc.codtypthr = vcom.codtypthr
       |            AND vloc.codthr = vcom.codthr
       |            LEFT JOIN TicketDocCrm tic ON tic.codvdt7 = vloc.codthr
       |            LEFT JOIN vcliff vc ON vc.codcli = tic.codcli
       |            LEFT JOIN $swamp_db.DBU_vclieml vcli ON vc.codcli = vcli.codcli
     """.stripMargin)

ORI_TktVCRAgency.registerTempTable("ORI_TktVCRAgency")

[...]

var TMP_workTemp = sqc.sql(
    s"""
       |SELECT *
       |FROM TicketDocCrm
       |            WHERE CODPNRREF != ''
       |            AND (DESRTGSTS LIKE '%USED%'
       |            OR DESRTGSTS LIKE '%OK%'
       |            OR DESRTGSTS LIKE '%CTRL%'
       |            OR DESRTGSTS LIKE '%RFND%'
       |            OR DESRTGSTS LIKE '%RPRT%'
       |            OR DESRTGSTS LIKE '%LFTD%'
       |            OR DESRTGSTS LIKE '%CKIN%')
     """.stripMargin)

TMP_workTemp.registerTempTable("TMP_workTemp")

var TMP_workTemp1 = sqc.sql(
    s"""
       |SELECT *
       |FROM TMP_workTemp w
       |INNER JOIN
       |    (SELECT CODTKTNUM as CODTKTNUM_T
       |    FROM (
       |        SELECT CODCLI, CODTKTNUM, COUNT(*) as n
       |        FROM TMP_workTemp
       |        GROUP BY CODCLI, CODTKTNUM
       |        HAVING n > 1)
       |    a) b
       |ON w.CODTKTNUM = b.CODTKTNUM_T
     """.stripMargin).drop("CODTKTNUM_T")

[...]

集群由 2 个 master 和 7 个 worker 组成。每个节点有:

  • 16核cpu
  • 110 GB 内存

Spark 在 YARN 上运行。

有没有人知道为什么在 Spark 中处理数据之前从原始文本转换为 Parquet 格式没有任何性能改进?

【问题讨论】:

    标签: scala apache-spark hive parquet snappy


    【解决方案1】:

    简短的回答。

    对于所有类型的查询,parquet 的性能都优于原始文本数据是不正确的。

    TLDR;

    Parquet 是 columnar store(描述什么是列式存储,表中的每一列都存储在单独的文件中,而不是存储行的文件中),这种模式(列式存储)提高了分析工作负载的性能(@ 987654322@)。

    我可以举一个例子说明为什么以柱状方式(如镶木地板)存储数据可能会显着提高查询性能。假设您有一个包含 300 列的表,并且您想要运行以下查询。

    SELECT avg(amount)
    FROM my_big_table 
    

    在上面的查询中,您只关心列金额的平均值。

    如果 spark 必须首先在原始文本上执行此操作,它将使用您提供的模式来拆分行,然后解析数量列,这需要相当多的计算时间来解析来自 300 个奇数列的数量列my_big_table。

    如果 spark 必须从 parquet 存储中获取平均金额,它必须只读取 amount-column-data 的 parquet 块(请记住,表的每一列都单独存储在 parquet )。 Parquet 可能会通过存储大量元数据和使用列级压缩来进一步提高性能。

    你应该阅读这个so post

    现在回到您的问题,您的大多数查询都在运行 SELECT *,这意味着您正在将所有数据读入 spark 中,然后加入或过滤一些值。在第二个查询中,使用 parquet 表的查询不会有太大的性能提升,因为您正在读取所有列,并且 parquet 在这里将是一个更昂贵的选择,因为您最终可能会读取更多您可能在原始文件中完成的文件-文本。

    在少数情况下,镶木地板的过滤速度更快,但并非总是如此,这取决于您的数据。

    总而言之,您应该根据要运行的查询类型和拥有的数据类型来选择数据存储。

    【讨论】:

    • 是的,我在想什么。使用 Avro 或 Sequencefile 等无列二进制格式可能会获得更好的性能,但我必须考虑转换的时间。谢谢你的好点。
    • 数据存储是一回事,您还应该考虑对数据集进行分区和分桶,以避免洗牌。如果您对答案感到满意,请接受答案。
    【解决方案2】:

    观察到的几个点:

    • hive.exec.compress.output=true -- 这将确保 Hive 查询的最终输出将被压缩。但在这种情况下,您是使用 Spark 从 Hive 读取数据,因此这不会对性能产生任何影响。
    • 检查数据帧的分区,确保有足够的数据帧分区,以便执行器并行处理数据。

    检查分区:

    vcliff.rdd.getNumPartitions
    
    • 按最常用的数据帧列对数据帧进行分区,以便 Spark 在执行连接等聚合时避免随机播放。如果最常用的列有更多不同的值,而不是分区,您可以在该列上使用 Bucketing,这样 Spark 将在分区之间均匀分布数据,而不是偏斜到一个或两个。

      vcliff.repartition([numPartitions], "codcli")

      TicketDocCrm.repartition([numPartitions], "DESRTGSTS")

    【讨论】:

    • a) 分区不会“避免洗牌” b) SQL 不会“失去 Spark Dataframes 提供的优化”
    猜你喜欢
    • 1970-01-01
    • 2018-02-04
    • 2020-05-17
    • 1970-01-01
    • 2017-03-27
    • 2017-07-26
    • 1970-01-01
    • 2016-11-04
    相关资源
    最近更新 更多