【问题标题】:How to access broadcasted DataFrame in Spark如何在 Spark 中访问广播的 DataFrame
【发布时间】:2018-01-09 19:55:26
【问题描述】:

我创建了两个数据框,它们来自 Hive 表(PC_ITM 和 ITEM_SELL)并且大小很大,我正在使用这些数据框 通过注册为表经常在SQL查询中。但是因为这些很大,所以需要很多时间 得到查询结果。所以我将它们保存为镶木地板文件,然后读取它们并注册为临时表。但我仍然没有获得良好的性能,所以我已经广播了这些数据帧,然后注册为如下表。

PC_ITM_DF=sqlContext.parquetFile("path")
val PC_ITM_BC=sc.broadcast(PC_ITM_DF)
val PC_ITM_DF1=PC_ITM_BC
PC_ITM_DF1.registerAsTempTable("PC_ITM")

ITM_SELL_DF=sqlContext.parquetFile("path")
val ITM_SELL_BC=sc.broadcast(ITM_SELL_DF)
val ITM_SELL_DF1=ITM_SELL_BC.value
ITM_SELL_DF1.registerAsTempTable(ITM_SELL)


sqlContext.sql("JOIN Query").show

但我仍然无法实现性能,因为它所花费的时间与未广播这些数据帧时的时间相同。

谁能判断这是否是广播和使用它的正确方法?`

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    您实际上并不需要“访问”广播数据帧 - 您只需使用它,Spark 就会在后台实现广播。 broadcast function 效果很好,并且比sc.broadcast 方法更有意义。

    如果一次评估所有内容,可能很难理解时间都花在了哪里。

    您可以将代码分解为多个步骤。此处的关键是执行一个操作并在您在加入中使用它们之前保存您想要广播的数据帧

    // load your dataframe
    PC_ITM_DF=sqlContext.parquetFile("path")
    
    // mark this dataframe to be stored in memory once evaluated
    PC_ITM_DF.persist()
    
    // mark this dataframe to be broadcast
    broadcast(PC_ITM_DF)
    
    // perform an action to force the evaluation
    PC_ITM_DF.count()
    

    这样做将确保数据帧是

    • 加载到内存中(持久)
    • 注册为临时表以用于您的 SQL 查询
    • 标记为广播,因此将发送给所有执行者

    当您现在运行 sqlContext.sql("JOIN Query").show 时,您现在应该在 Spark UI 的 SQL 选项卡中看到一个“广播哈希连接”。

    【讨论】:

    • 广播 RDD 有什么好处? RDD 代表弹性分布式数据集。广播消除了 RDD 的分布式特性。我可以看到您将数据从 RDD 收集到内存并广播的用例。我什至不相信这是可能的。如果您查看此article,它会说..“要广播 RDD,您需要首先在驱动程序节点上收集()它。”你曾经在实践或测试中使用过这个吗?
    • @AlexNaspo 是的,我一直在使用它。好处是数据在所有节点上都是完全可用的——它不再是分布式的——这有助于提高加入时的性能。例如,考虑一个包含美国每个人及其邮政编码的 DataFrame,然后是一个包含邮政编码 -> 州的表格。加入这些需要大量的洗牌。将相对较小的 zip->state 数据帧广播到所有节点,无需重新洗牌。
    • 您正在广播保存在内存中的数据帧,而不是分发的。那是对的吗? Spark 建议在您的数据中添加一个分区器,以减少加入时的 shuffle 量。 @kirkbroadhurt
    • 我告诉 Spark 应该广播数据帧。它保持“分布式”,直到需要它(例如用于连接),此时 Spark 的 Catalyst 优化器知道我希望它向每个节点发送数据帧的完整副本。
    • 好吧,我误会了。这是在数据帧适合内存的条件下。在您的邮政编码示例中,效果很好。在分布式数据帧大于您的内存的情况下,广播似乎不是正确的方法。在这种情况下,将数据持久化和分区相结合将是一种适用于任何大小数据的解决方案。
    【解决方案2】:

    我会将 rdds 缓存在内存中。下次需要它们时,spark 将从内存中读取 RDD,而不是每次都从头开始生成 RDD。这是快速入门docs 的链接。

    val PC_ITM_DF = sqlContext.parquetFile("path")
    PC_ITM_DF.cache()
    PC_ITM_DF.registerAsTempTable("PC_ITM")
    
    val ITM_SELL_DF=sqlContext.parquetFile("path")
    ITM_SELL_DF.cache()
    ITM_SELL_DF.registerAsTempTable("ITM_SELL")
    sqlContext.sql("JOIN Query").show
    

    rdd.cache() 是rdd.persist(StorageLevel.MEMORY_ONLY) 的简写。您可以选择几个级别的持久性,以防您的数据太大而无法仅用于内存持久性。这里是list of persistence options.,如果你想手动从缓存中删除RDD,你可以调用rdd.unpersist()

    如果您更喜欢广播数据。您必须先在驱动程序上收集它,然后才能广播它。这要求您的 RDD 适合您的驱动程序(和执行程序)的内存。

    【讨论】:

    • 这并没有回答最初的问题,即如何广播 DataFrame。仅当您多次加载它(即重用)时,持久化才会有所帮助。加入两个分布式数据集时没有帮助。
    • @KirkBroadhurst 他说数据很大,使用频率很高
    • @AlexNaspo 稍微偏离了原来的问题:RDD fits in memory 因此这意味着我不能广播数据,直到我可以将它收集到驱动程序的主存储器中?我通常使用自己的笔记本电脑作为大型集群上的驱动程序和主/从。那么这是我可能很快面临的限制吗?
    • 根据@KirkBroadhurst,您可以广播RDD,执行器将在需要时收集数据。
    【解决方案3】:

    此时您无法在 SQL 查询中访问广播数据帧。您只能通过数据帧使用广播数据帧。

    参考:https://issues.apache.org/jira/browse/SPARK-16475

    【讨论】:

    • 目前的解决方案是先在dataframe api中广播df或表,将broadcast函数的返回值注册为临时表,然后在SQL查询中调用该临时表。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-04-02
    • 2017-01-26
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多