【问题标题】:Understanding spark physical plan了解火花物理计划
【发布时间】:2016-09-27 02:26:02
【问题描述】:

我正在尝试了解 spark 的物理计划,但我不了解某些部分,因为它们似乎与传统的 rdbms 不同。例如,在下面的这个计划中,它是一个关于对 hive 表进行查询的计划。查询是这样的:

select
        l_returnflag,
        l_linestatus,
        sum(l_quantity) as sum_qty,
        sum(l_extendedprice) as sum_base_price,
        sum(l_extendedprice * (1 - l_discount)) as sum_disc_price,
        sum(l_extendedprice * (1 - l_discount) * (1 + l_tax)) as sum_charge,
        avg(l_quantity) as avg_qty,
        avg(l_extendedprice) as avg_price,
        avg(l_discount) as avg_disc,
        count(*) as count_order
    from
        lineitem
    where
        l_shipdate <= '1998-09-16'
    group by
        l_returnflag,
        l_linestatus
    order by
        l_returnflag,
        l_linestatus;


== Physical Plan ==
Sort [l_returnflag#35 ASC,l_linestatus#36 ASC], true, 0
+- ConvertToUnsafe
   +- Exchange rangepartitioning(l_returnflag#35 ASC,l_linestatus#36 ASC,200), None
      +- ConvertToSafe
         +- TungstenAggregate(key=[l_returnflag#35,l_linestatus#36], functions=[(sum(l_quantity#31),mode=Final,isDistinct=false),(sum(l_extendedpr#32),mode=Final,isDistinct=false),(sum((l_extendedprice#32 * (1.0 - l_discount#33))),mode=Final,isDistinct=false),(sum(((l_extendedprice#32 * (1.0l_discount#33)) * (1.0 + l_tax#34))),mode=Final,isDistinct=false),(avg(l_quantity#31),mode=Final,isDistinct=false),(avg(l_extendedprice#32),mode=Fl,isDistinct=false),(avg(l_discount#33),mode=Final,isDistinct=false),(count(1),mode=Final,isDistinct=false)], output=[l_returnflag#35,l_linestatus,sum_qty#0,sum_base_price#1,sum_disc_price#2,sum_charge#3,avg_qty#4,avg_price#5,avg_disc#6,count_order#7L])
            +- TungstenExchange hashpartitioning(l_returnflag#35,l_linestatus#36,200), None
               +- TungstenAggregate(key=[l_returnflag#35,l_linestatus#36], functions=[(sum(l_quantity#31),mode=Partial,isDistinct=false),(sum(l_exdedprice#32),mode=Partial,isDistinct=false),(sum((l_extendedprice#32 * (1.0 - l_discount#33))),mode=Partial,isDistinct=false),(sum(((l_extendedpri32 * (1.0 - l_discount#33)) * (1.0 + l_tax#34))),mode=Partial,isDistinct=false),(avg(l_quantity#31),mode=Partial,isDistinct=false),(avg(l_extendedce#32),mode=Partial,isDistinct=false),(avg(l_discount#33),mode=Partial,isDistinct=false),(count(1),mode=Partial,isDistinct=false)], output=[l_retulag#35,l_linestatus#36,sum#64,sum#65,sum#66,sum#67,sum#68,count#69L,sum#70,count#71L,sum#72,count#73L,count#74L])
                  +- Project [l_discount#33,l_linestatus#36,l_tax#34,l_quantity#31,l_extendedprice#32,l_returnflag#35]
                     +- Filter (l_shipdate#37 <= 1998-09-16)
                        +- HiveTableScan [l_discount#33,l_linestatus#36,l_tax#34,l_quantity#31,l_extendedprice#32,l_shipdate#37,l_returnflag#35], astoreRelation default, lineitem, None

我对计划的理解是:

  1. 首先从 Hive 表扫描开始

  2. 然后使用where条件进行过滤

  3. 然后project得到我们想要的列

  4. 然后是钨骨料?

  5. 然后是钨交换?

  6. 然后再次聚合钨?

  7. 然后 ConvertToSafe?

  8. 然后对最终结果进行排序

但我不了解 4、5、6 和 7 步骤。你知道它们是什么吗?我正在寻找有关这方面的信息,以便了解该计划,但我没有找到任何具体的信息。

【问题讨论】:

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


    【解决方案1】:

    让我们看看您使用的 SQL 查询的结构:

    SELECT
        ...  -- not aggregated columns  #1
        ...  -- aggregated columns      #2
    FROM
        ...                          -- #3
    WHERE
        ...                          -- #4
    GROUP BY
        ...                          -- #5
    ORDER BY
        ...                          -- #6
    

    正如你已经怀疑的那样:

    • Filter (...) 对应于WHERE 子句中的谓词(#4)
    • Project ... 将列数限制为(#1 和 #2,以及 #4 / #6 如果在 SELECT 中不存在)所要求的列数
    • HiveTableScan 对应于FROM 子句(#3)

    其余部分可归为如下:

    • #2 来自 SELECT 子句 - functions TungstenAggregates 中的字段
    • GROUP BY 子句 (#5):

      • TungstenExchange/哈希分区
      • keyTungstenAggregates 中的字段
    • #6 - ORDER BY 子句。

    Project Tungsten 总体上描述了 Spark DataFrames (-sets) 使用的一组优化,包括:

    • 使用sun.misc.Unsafe 进行显式内存管理。它意味着“本机”(堆外)内存使用和显式内存分配/在 GC 管理之外释放。这些转换对应于执行计划中的ConvertToUnsafe / ConvertToSafe 步骤。你可以从Understanding sun.misc.Unsafe 了解一些关于 unsafe 的有趣细节
    • 代码生成 - 不同的元编程技巧旨在生成在编译期间得到更好优化的代码。你可以把它想象成一个内部的 Spark 编译器,它可以将漂亮的函数式代码重写为丑陋的 for 循环。

    您可以从Project Tungsten: Bringing Apache Spark Closer to Bare Metal 了解有关钨的更多信息。 Apache Spark 2.0: Faster, Easier, and Smarter 提供了一些代码生成示例。

    TungstenAggregate 出现两次,因为数据首先在每个分区上本地聚合,然后是混洗,最后合并。如果你熟悉 RDD API,这个过程大致相当于reduceByKey。

    如果执行计划不清楚,您也可以尝试将生成的DataFrame 转换为RDD 并分析toDebugString 的输出。

    【讨论】:

    • 感谢您的回答。我只是没有清楚地理解这部分“SELECT 子句中的#2 - TungstenAggregates 中的函数字段”。如果你能解释得更好就更好了!
    • Functions 字段列出了在给定阶段执行的所有聚合,而Key 字段描述了分组。它是df.groupBy(*key).agg(*functions)。
    • @zero323 GROUP BY clause (#4): 应该是 GROUP BY clause (#5):。我想自己编辑,但 stackOverflow 告诉我 Edits must be at least 6 characters; is there something else to improve in this post?
    【解决方案2】:

    Tungsten 是 Spark 自 1.4 以来的新内存引擎,它管理 JVM 外部的数据以节省一些 GC 开销。您可以想象这样做涉及将数据从 JVM 复制到 JVM。而已。在 Spark 1.5 中,您可以通过 spark.sql.tungsten.enabled 关闭 Tungsten,然后您会看到“旧”计划,在 Spark 1.6 中,我认为您不能再关闭它了。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2020-05-26
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-03-10
      相关资源
      最近更新 更多