【问题标题】:Alternative of hive join蜂巢连接的替代方案
【发布时间】:2017-02-20 13:19:48
【问题描述】:

我的蜂巢中有两个视图

+------------+
| Table_1    |
+------------+
| hash       |
| campaignId |
+------------+

+-----------------+
| Table_2         |
+-----------------+
| campaignId      |
| accountId       |
| parentAccountID |
+-----------------+

现在我必须获取按 accountId 和 parentAccountID 过滤的“Table_1”数据,为此我编写了以下查询:

SELECT /*+ MAPJOIN(T2) */ T1.hash, COUNT(T1.campaignId) num_campaigns
FROM Table_1 T1
JOIN Table_2 T2 ON T1.campaignId = T2.campaignId
WHERE (T2.accountId IN ('aid1', 'aid2') OR T2.parentAccountID IN ('aid1', 'aid2')
GROUP BY T1.hash

此查询正在运行,但速度很慢。有没有更好的替代方法(JOIN)?

我正在通过 spark 从 kafka 读取 Table_1。
幻灯片持续时间为 5 秒
窗口持续时间为 2 分钟

虽然 Table_2 在 RDBMS 中,但 spark 正在通过 jdbc 读取,并且有 4500 条记录。

每 5 秒,kafka 以 CSV 格式抽取大约 2K 条记录。
我需要在 5 秒内处理数据,但目前需要 8 到 16 秒。

根据建议:

  1. 我已分别按列 campaignId 和哈希对 Table_1 进行了重新分区。
  2. 我已分别按 accountId 和 parentAccountID 列对 Table_2 进行了重新分区。
  3. 我已经实现了 MAPJOIN。

但仍然没有改善。

注意:如果我删除窗口持续时间,那么该过程确实会在一段时间内执行。可能是因为要处理的数据较少。但这不是要求。

【问题讨论】:

  • (1) 您如何定义“慢”? (2) 我们在谈论什么卷?
  • 我在 Spark Streaming 中使用它,每 5 秒处理一次数据。但它需要超过 10 秒来处理每批。 (注意:幻灯片持续时间为 5 秒,但窗口持续时间为 10 秒)。每批大约有 24K 条记录。
  • 你需要'hash'列吗?
  • 是的。我按它分组。
  • 您是否对每个流上的每个 RDD(DStream)进行此查询?因为我觉得在每个流中创建 HiveContext 和 Hive 查询可能很昂贵。

标签: sql hadoop hive query-optimization hiveql


【解决方案1】:

使用正确的索引,以下操作会更快:

SELECT T1.*
FROM Table_1 T1 JOIN
     Table_2 T2
     ON T1.campaignId = T2.campaignId
WHERE T2.accountId IN ('aid1', 'aid2') 
UNION ALL
SELECT T1.*
FROM Table_1 T1 JOIN
     Table_2 T2
     ON T1.campaignId = T2.campaignId
WHERE T2.parentAccountID IN ('aid1', 'aid2') AND
      T2.accountId NOT IN ('aid1', 'aid2') ;

第一个可以考虑Table_2(accountId, campaignId) 上的索引,第二个可以考虑Table_2(parentAccountID, accountId, campaignId) 上的索引。

【讨论】:

    【解决方案2】:

    由于这是我们谈论的 Hive,您需要了解的不仅仅是传统 DBMS。

    • 减少 IO。为您的数据使用压缩的列格式。 ORC 或镶木地板。 不是 RC。首先执行此操作,将您的表转换为 ORC。除非数据被压缩并呈柱状,否则其他任何事情都不会产生太大影响。
    • 为 Hive 选择正确的 JOIN 策略。这个old 2011 paper 仍然有用。
    • Bucketize你的桌子
    • 使用现代执行引擎:Tez 或 Spark。

    【讨论】:

    • 感谢@remus 的快速回复。我正在从 kafka 流中读取数据。数据是 csv 格式,我无法更改。我正在使用火花发动机。
    • CSV 你在没有桨的情况下上岸...检查你得到的 JOIN 风格。如果其中一个表足够小,也许您可​​以使用 Map 连接(Table_1 可能?)
    • 顺便说一句,如果数据是流式传输的,您应该分区/分桶并添加时间过滤器。这应该会大大减少扫描大量数据的需要。
    • 另外,我建议你阅读Streaming Data Ingest
    • 我已经实现了分区和 Map Join,但没有发现任何改进。 :( 我还通过原始问题更新了更多详细信息。
    【解决方案3】:

    如果过滤后的 T2 足够小以适合内存,请尝试重写查询并将过滤器移动到子查询中,看看是否会在映射器上执行连接。此外,您不需要 T2 中的列,可以使用左半连接代替内连接:

    set hive.cbo.enable=true; 
    set hive.auto.convert.join=true;
    
    SELECT T1.* 
    FROM Table_1 T1 
    LEFT SEMI JOIN 
         (select campaignId  from Table_2 T2 
            where T2.accountId IN ('aid1', 'aid2') 
               OR T2.parentAccountID IN ('aid1', 'aid2')
         ) T2 ON T1.campaignId = T2.campaignId 
    ;
    

    【讨论】:

      【解决方案4】:

      我建议您使用本机 Spark 转换而不是 HiveSQL:

      1.将Table_2 (RDBMS)中的数据读入RDD并放入缓存 例如:

      rddTbl1.map(campaignIdKey, (accountId, parentAccountId)) //filter out before getting into RDD if needed
      rddTbl2.cache()
      

      2.现在读取Table_1流(Kafka)

      //get campaigns of relevant account & parentaccountid
      val rddTbl2_1 = rddTbl2.filter(x => x._2._1.equals("aid1") || x._2._1.equals("aid2") || x._2._2.equals("aid1") || x._2._2.equals("aid2"))
      
      dstream.foreachRDD{ rddTbl1 =>
        rddTbl1.map(x => x._2.split(",")).
                map(x => (x(1), x(2)). //campaignId, hash
                join(rddTbl2_1).
                map(x => (x._2._1, 1)). //get (hash,1)
                reduceByKey(_+_).
                foreach(println) //save it if needed
      }
      

      【讨论】:

        【解决方案5】:

        好的..

        这是我最终所做的。

        我创建了 Table_2 的哈希。
        然后通过使用广播变量,我将该数据传递给每个节点。

        这样就省去了加入的麻烦。

        感谢大家的时间和帮助。 快乐编码:)

        【讨论】:

          猜你喜欢
          • 1970-01-01
          • 1970-01-01
          • 2015-05-15
          • 1970-01-01
          • 1970-01-01
          • 2021-04-21
          • 1970-01-01
          • 2014-10-13
          • 2018-06-08
          相关资源
          最近更新 更多