【发布时间】: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