【发布时间】:2018-05-07 06:52:32
【问题描述】:
我正在处理一个文本文件并将转换后的行从 Spark 应用程序写入弹性搜索,如下所示
input.write.format("org.elasticsearch.spark.sql")
.mode(SaveMode.Append)
.option("es.resource", "{date}/" + dir).save()
这运行速度非常慢,写入 287.9 MB / 1513789 条记录大约需要 8 分钟。
鉴于始终存在网络延迟,我如何调整 spark 和 elasticsearch 设置以使其更快。
我在本地模式下使用 spark,有 16 个内核和 64GB RAM。 我的 elasticsearch 集群有 1 个主节点和 3 个数据节点,每个节点有 16 个内核和 64GB。
我正在阅读如下文本文件
val readOptions: Map[String, String] = Map("ignoreLeadingWhiteSpace" -> "true",
"ignoreTrailingWhiteSpace" -> "true",
"inferSchema" -> "false",
"header" -> "false",
"delimiter" -> "\t",
"comment" -> "#",
"mode" -> "PERMISSIVE")
....
val input = sqlContext.read.options(readOptions).csv(inputFile.getAbsolutePath)
【问题讨论】:
-
输入中有多少个分区? spark 是否与 elasticsearch 共享资源?您可以随时更改默认写入批量大小,我相信默认为 500。
-
实际上我正在使用 csv 插件
val input = sqlContext.read.options(readOptions).csv(inputFile.getAbsolutePath)创建数据框。我不知道它创建了多少个分区。 Spark 不与 elasticsearch 共享资源。 -
你的文件有多大?它是分区的还是一个大集团?
-
它是 70MB gz 文件,包含网络访问日志。这是一个未分区的文件
-
我想你在你的 readOptions 中推断Schema,这将导致扫描数据两次。你能打印 df.rdd.getNumPartitions 的输出吗?
标签: apache-spark elasticsearch elasticsearch-5 elasticsearch-spark