【问题标题】:spark cross join memory leak火花交叉连接内存泄漏
【发布时间】:2018-03-04 06:44:45
【问题描述】:

我有两个表要交叉连接,

表 1:查询 300M 行
表2:产品描述3000行

以下查询进行交叉连接并计算元组之间的分数,并选择前 3 个匹配项,

query_df.repartition(10000).registerTempTable('queries')

product_df.coalesce(1).registerTempTable('products')

CREATE TABLE matches AS
SELECT *
FROM
  (SELECT *,
          row_number() over (partition BY a.query_id
                             ORDER BY 0.40 + 0.15*score_a + 0.20*score_b + 0.5*score_c DESC) AS rank
   FROM
     (SELECT /*+ MAPJOIN(b) */ a.query_id,
                               b.product_id,
                               func_a(a.qvec,b.pvec) AS score_a,
                               func_b(a.qvec,b.pvec) AS score_b,
                               func_c(a.qvec,b.pvec) AS score_c
      FROM queries a CROSS
      JOIN products b) a) a
WHERE rn <= 3

我的火花集群如下所示,

MASTER="yarn-client" /opt/mapr/spark/spark-1.6.1/bin/pyspark --num-executors 22 --executor-memory 30g --executor-cores 7 --driver-memory 10g --conf spark.yarn.executor.memoryOverhead=10000 --conf spark.akka.frameSize=2047

现在的问题是,正如预期的那样,由于内存泄漏,由于产生了极大的临时数据,作业在几个阶段后失败。我正在寻找一些帮助/建议来优化上述操作,以使作业应该能够在选择下一个 query_id 之前运行 query_id 的匹配和过滤操作,在一种并行方式 - 类似于针对查询表的 for 循环中的排序。如果工作缓慢但成功,我可以接受,因为我可以请求更大的集群。

上述查询适用于较小的查询表,比如有 10000 条记录的表。

【问题讨论】:

    标签: apache-spark pyspark apache-spark-sql spark-dataframe


    【解决方案1】:

    在您想要连接表 A(大)和表 B(小)的场景中,最佳做法是利用 广播连接

    https://stackoverflow.com/a/39404486/1203837 中给出了清晰的概述。

    希望这会有所帮助。

    【讨论】:

    • 我已经使用 Hive 语法 /*+ MAPJOIN(b) */ 使用广播连接,它包含在查询中。
    【解决方案2】:

    Spark 中的笛卡尔连接或交叉连接非常昂贵。我建议使用内部连接加入表并首先保存输出数据。然后使用该数据框进行进一步聚合。

    如果较小的表不够小,映射连接或广播连接有时可能会失败的一个小建议。除非您确定小表的大小,否则请不要使用广播连接。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-03-16
      • 1970-01-01
      • 1970-01-01
      • 2018-11-24
      相关资源
      最近更新 更多