【问题标题】:Apache Spark Map Job SlowApache Spark Map 作业速度慢
【发布时间】:2017-12-19 03:10:46
【问题描述】:

我一直在试验 Apache Spark,看看它是否可用于为我们存储在 Elasticsearch 集群中的数据创建分析引擎。我发现对于任何重要的 RDD 大小(即几百万条记录),即使是最简单的操作也需要一分钟以上的时间。

比如我做了这个简单的测试程序:

package es_spark;

import java.util.Map;

import org.apache.spark.SparkConf;
import org.apache.spark.api.java.JavaPairRDD;
import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.api.java.JavaSparkContext;
import org.elasticsearch.spark.rdd.api.java.JavaEsSpark;

public class Main {

    public static void main (String[] pArgs) {

        SparkConf conf = new SparkConf().setAppName("Simple Application");
        conf.set("es.nodes", pArgs[0]);

        JavaSparkContext sc = new JavaSparkContext(conf);

        long start = System.currentTimeMillis();
        JavaPairRDD<String, Map<String, Object>> esRDD = JavaEsSpark.esRDD(sc, "test3");
        long numES = esRDD.count();
        long loadStop = System.currentTimeMillis();

        JavaRDD<Integer> dummyRDD = esRDD.map(pair -> {return 1;});
        long numDummy = dummyRDD.count();
        long mapStop = System.currentTimeMillis();

        System.out.println("ES Count: " + numES);
        System.out.println("ES Partitions: " + esRDD.getNumPartitions());

        System.out.println("Dummy Count: " + numDummy);
        System.out.println("Dummy Partitions: " + dummyRDD.getNumPartitions());

        System.out.println("Data Load Took: " + (loadStop - start) + "ms");
        System.out.println("Dummy Map Took: " + (mapStop - loadStop) + "ms");

        sc.stop();
        sc.close();
    }
}

我在一个有 3 个从属设备的 spark 集群上运行它,每个从属设备有 14 个内核和 49.0GB 的 RAM。使用以下命令:

./bin/spark-submit --class es_spark.Main --master spark://<master_ip>:7077 ~/es_spark-0.0.1.jar <elasticsearch_main_ip>

输出是:

ES Count: 8140270
ES Partitions: 80
Dummy Count: 8140270
Dummy Partitions: 80
Data Load Took: 108059ms
Dummy Map Took: 104128ms

对 8+ 百万条记录执行虚拟映射作业需要 1.5+ 分钟。鉴于地图作业什么都不做,我发现这种性能出奇地低。我做错了什么还是这与 Spark 的正常性能有关?

我也尝试过调整--executor-memory--executor-cores,但没有太大区别。

【问题讨论】:

    标签: java apache-spark elasticsearch


    【解决方案1】:

    鉴于地图作业不执行任何操作,您会发现此性能低得惊人。

    地图作业什么也不做。它必须从弹性搜索中获取完整的数据集。因为数据没有被缓存,所以它会发生两次,每个动作一次。这个时间还包括一些初始化时间。

    您衡量的总体情况:

    • ES 查询的时间。
    • Spark 集群和 ES 之间的网络延迟。

    还有一些次要的东西,例如:

    • 执行程序 JVM 的完全初始化时间。
    • 可能是 GC 暂停时间。

    【讨论】:

    • 我还尝试在计数后立即在 esRDD 上调用缓存:long numES = esRDD.count(); esRDD.cache();这不会阻止来自 Elastic Search 的双重提取吗?地图的时间变化不大。
    • 我将上面的示例代码更改为在 esRDD.count() 之前调用 esRDD.cache()。不幸的是,它对地图性能没有太大帮助。虚拟地图耗时:113298ms
    • 不过,您可能正在做一些事情,因为我计算了 esRDD 的数量,这使得 整个 测试程序在大约 1.5 分钟内运行。所以地图的运行时间几乎没有。
    • user8371915 是正确的,映射操作至少返回到弹性搜索的某些分区。添加 esRDD.persist(StorageLevel.MEMORY_AND_DISK());显着减少了虚拟地图操作的运行时间。
    【解决方案2】:

    通常,除非您看到 OOM 失败或严重的 GC 或溢出到磁盘作为瓶颈,否则不值得更改执行程序内存。当你改变它时,你也应该减少 spark.memory.fraction。对于您的工作,这不太可能有帮助。

    Spark 存在启动成本,这使得它对于较小的数据负载效率相对较低。您应该能够将您的启动优化到不到一分钟,但对于非常大的批量负载,而不是实时分析,它仍然更实用。

    我建议您使用 DataFrame API 而不是 RDD。对于上面的简单示例操作,这无关紧要,但随着事情变得更加复杂,您更有可能从性能优化中受益。

    例如sql.read.format("es").load("test3")

    要排查导致运行缓慢的原因,您可以查看 Spark UI。你真的得到了并行性吗?所有作业是否在大致相同的时间内执行?另一个可能导致缓慢的原因是您的集群和 ES 服务器之间的网络问题。

    【讨论】:

    • 我并不特别关心 Elasticsearch 的加载时间。我们可以忽略加载时间信息。我更担心加载后执行无操作映射作业需要 1.5 分钟。
    猜你喜欢
    • 1970-01-01
    • 2017-05-30
    • 1970-01-01
    • 2019-11-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多