【问题标题】:Spark: best practice on memory-heavy join operationsSpark:重内存连接操作的最佳实践
【发布时间】:2019-01-14 14:37:17
【问题描述】:

我有一个 spark 程序,它涉及对大型 Hive 表(数百万行和数百列)的连接操作。 这些连接期间使用的内存非常高。我想了解在 Spark on YARN 中处理这种情况的最佳方法,以使作业成功完成而不会出现内存错误。该集群由 7 个工作人员组成,每个工作人员具有 110 GB 的内存和 16 个内核。 考虑以下 scala 代码:

object Model1Prep {

    def main(args: Array[String]): Unit = {

        val conf = new SparkConf().setAppName("Modello1_Spark")
        conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
        conf.set("spark.io.compression.codec", "org.apache.spark.io.LZFCompressionCodec")
        val sc = new SparkContext(conf)
        val hc = new HiveContext(sc)
        import hc.implicits._


        hc.sql("SET hive.exec.compress.output=true")
        hc.sql("SET parquet.compression=SNAPPY")
        hc.sql("SET spark.sql.parquet.compression.codec=snappy")


        // loading tables on dataframes
        var tableA = hc.read.table("TA")
        var tableB = hc.read.table("TB")
        var tableC = hc.read.table("TC")
        var tableD = hc.read.table("TD")


        // registering tables
        tableA.registerTempTable("TA")
        tableB.registerTempTable("TB")
        tableC.registerTempTable("TC")
        tableD.registerTempTable("TD")


        var join1 = hc.sql("""
            SELECT 
                [many fields]
            FROM TA a 
            JOIN TB b ON a.field = b.field      
            LEFT JOIN TC c ON a.field = c.field         
            WHERE [conditions]
        """)


        var join2 = hc.sql("""
            SELECT 
                [many fields]
            FROM TA a 
            LEFT JOIN TD d ON a.field = d.field
            WHERE [conditions]
        """)


        // [other operations]


        sc.close()
    }
}

考虑到连接操作在内存上真的很昂贵,我最好的选择是什么? 我知道数据帧可以同时保存在内存和磁盘上,可能使用序列化在内存中更紧凑,但代价是更慢的反序列化处理时间(更多关于herehere)。 从上面的代码中,表 TA 在两个连接中都使用了,因此持久化它是有意义的:

    //[...]        

    // persisting
    tableA.persist(StorageLevel.MEMORY_AND_DISK_SER_2)

    // registering tables
    tableA.registerTempTable("TA")
    tableB.registerTempTable("TB")
    tableC.registerTempTable("TC")
    tableD.registerTempTable("TD")

    //[...]

我也应该以同样的方式持久化其他表吗?还是有其他的东西可以让这段代码顺利运行并完成?

【问题讨论】:

    标签: scala apache-spark hadoop pyspark hadoop-yarn


    【解决方案1】:

    如果您知道要加入哪个字段并且它始终是同一个字段,那么 as this SO answer suggests,对已加入的表使用相同的分区器。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-11-30
      • 2021-05-17
      • 1970-01-01
      • 2015-03-25
      • 1970-01-01
      • 2012-09-01
      • 1970-01-01
      • 2015-08-03
      相关资源
      最近更新 更多