【发布时间】: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