【问题标题】:How to optimize a join?如何优化连接?
【发布时间】:2017-10-27 08:22:41
【问题描述】:

我有一个要加入表格的查询。如何优化以更快地运行它?

val q = """
          | select a.value as viewedid,b.other as otherids
          | from bm.distinct_viewed_2610 a, bm.tets_2610 b
          | where FIND_IN_SET(a.value, b.other) != 0 and a.value in (
          |   select value from bm.distinct_viewed_2610)
          |""".stripMargin
val rows = hiveCtx.sql(q).repartition(100)

表格说明:

hive> desc distinct_viewed_2610;
OK
value                   string

hive> desc tets_2610;
OK
id                      int                                         
other                   string 

数据如下所示:

hive> select * from distinct_viewed_2610 limit 5;
OK
1033346511
1033419148
1033641547
1033663265
1033830989

hive> select * from tets_2610 limit 2;
OK

1033759023
103973207,1013425393,1013812066,1014099507,1014295173,1014432476,1014620707,1014710175,1014776981,1014817307,1023740250,1031023907,1031188043,1031445197

distinct_viewed_2610 表有 110 万条记录,我正在尝试通过拆分第二列从 tets_2610 表中获取类似的 id,该表有 200 000 行。

对于 100 000 条记录,使用两台机器完成这项工作需要 8.5 小时 一个具有 16 GB RAM 和 16 个内核 其次是 8 GB 内存和 8 个内核。

有没有办法优化查询?

【问题讨论】:

  • 对表进行分区(如果可能),使用 ORC 或 Parquet 以获得更好的文件格式,调整内存设置以便将较小的连接表放入内存(如果可能),等等...通常如何压缩 Spark 性能。问题是:您正在运行多少个执行程序?你给他们什么设置?您的数据是否存在偏差,导致一位执行者需要很长时间?
  • 感谢您的回复..需要澄清以下点..1。我从少数文档中读到,spark 不会检查数据以哪种格式存储它直接从路径访问数据。我对 ORC 和非 ORC 表的查询证明响应时间似乎相同?跨度>
  • 你在哪里读到的?当然 Spark 关心数据格式。如果您制作数据集,那么列式数据格式将比纯文本文件最佳。
  • 我有 4 个执行程序,并且刚刚附上了我的执行程序选项卡的快照,其中提供了有关内存分配的信息
  • 嗯,乍一看,您在较大的机器上只使用了 12 个内核,而在另一台机器上使用了 4 个内核。我猜4个核心被分配给别的东西。如果你想要更多的执行者,你可以设置 1 或 2 个核心。而且它看起来根本没有真正使用内存。

标签: apache-spark hive apache-spark-sql query-optimization


【解决方案1】:

现在你正在做笛卡尔连接。笛卡尔连接为您提供 1.1M*200K = 2200 亿行。笛卡尔加入后由where FIND_IN_SET(a.value, b.other) != 0过滤

分析您的数据。 如果“其他”字符串平均包含 10 个元素,那么将其分解将为您在表 b 中提供 220 万行。如果假设只有 1/10 的行会加入,那么由于 INNER JOIN,您将拥有 2.2M/10=220K 行。

如果这些假设是正确的,那么爆炸数组和连接将比笛卡尔连接+过滤器执行得更好。

select distinct a.value as viewedid, b.otherids
  from bm.distinct_viewed_2610 a
       inner join (select e.otherid, b.other as otherids 
                     from bm.tets_2610 b
                          lateral view explode (split(b.other ,',')) e as otherid
                  )b on a.value=b.otherid

你不需要这个:

and a.value in (select value from bm.distinct_viewed_2610)

对不起,我无法测试查询,请自己做。

【讨论】:

    【解决方案2】:

    如果您根据您的数据使用 orc 甲酸盐更改为镶木地板,我会说选择范围分区。

    选择适当的并行化以快速执行。

    我已经回答了下面的链接可能对你有帮助。

    Spark doing exchange of partitions already correctly distributed

    请仔细阅读

    http://dev.sortable.com/spark-repartition/

    【讨论】:

      猜你喜欢
      • 2011-06-05
      • 2022-01-24
      • 1970-01-01
      • 2012-11-03
      • 2011-04-29
      • 1970-01-01
      • 2011-05-10
      • 2018-02-26
      • 2011-05-22
      相关资源
      最近更新 更多