【问题标题】:Spark with Kafka streaming save to Elastic search slow performanceSpark 与 Kafka 流保存到 Elastic 搜索性能缓慢
【发布时间】:2018-08-23 18:48:34
【问题描述】:

我有一个数据列表,值基本上是一个 bson 文档(想想 json),每个 json 的大小从 5k 到 20k 不等。既可以是bson对象格式,也可以直接转成json:

Key, Value
--------
K1, JSON1
K1, JSON2
K2, JSON3
K2, JSON4

我希望 groupByKey 会产生:

K1, (JSON1, JSON2)
K2, (JSON3, JSON4)

所以当我这样做时:

val data = [...].map(x => (x.Key, x.Value))
val groupedData = data.groupByKey()
groupedData.foreachRDD { rdd =>
   //the elements in the rdd here are not really grouped by the Key
}

我对 RDD 的行为感到非常困惑。我在网上看了很多文章,包括来自 Spark 的官方网站:https://spark.apache.org/docs/0.9.1/scala-programming-guide.html

仍然无法实现我想要的。

-------- 已更新 ---------

基本上我确实需要按key进行分组,key就是要在Elasticsearch中使用的索引,这样我就可以通过Elasticsearch for Hadoop根据key进行批处理:

EsSpark.saveToEs(rdd);

我不能按分区做,因为 Elasticsearch 只接受 RDD。我尝试使用 sc.MakeRDD 或 sc.parallize,都告诉我它不可序列化。

我尝试使用:

EsSpark.saveToEs(rdd, Map(
          "es.resource.write" -> "{TheKeyFromTheObjectAbove}",
          "es.batch.size.bytes" -> "5000000")

配置文档在这里:https://www.elastic.co/guide/en/elasticsearch/hadoop/current/configuration.html

但是与不使用配置根据单个文档的值定义动态索引相比非常慢,我怀疑它正在解析每个 json 以动态获取值。

【问题讨论】:

标签: scala apache-spark elasticsearch


【解决方案1】:

这是一个例子。

import org.apache.spark.SparkConf
import org.apache.spark.rdd.RDD
import org.apache.spark.sql.SparkSession

object Test extends App {

  val session: SparkSession = SparkSession
    .builder.appName("Example")
    .config(new SparkConf().setMaster("local[*]"))
    .getOrCreate()
  val sc = session.sparkContext

  import session.implicits._

  case class Message(key: String, value: String)

  val input: Seq[Message] =
    Seq(Message("K1", "foo1"),
      Message("K1", "foo2"),
      Message("K2", "foo3"),
      Message("K2", "foo4"))

  val inputRdd: RDD[Message] = sc.parallelize(input)

  val intermediate: RDD[(String, String)] =
    inputRdd.map(x => (x.key, x.value))
  intermediate.toDF().show()
  //  +---+----+
  //  | _1|  _2|
  //  +---+----+
  //  | K1|foo1|
  //  | K1|foo2|
  //  | K2|foo3|
  //  | K2|foo4|
  //  +---+----+

  val output: RDD[(String, List[String])] =
    intermediate.groupByKey().map(x => (x._1, x._2.toList))
  output.toDF().show()
  //  +---+------------+
  //  | _1|          _2|
  //  +---+------------+
  //  | K1|[foo1, foo2]|
  //  | K2|[foo3, foo4]|
  //  +---+------------+

  output.foreachPartition(rdd => if (rdd.nonEmpty) {
    println(rdd.toList)
  })
  //  List((K1,List(foo1, foo2)))
  //  List((K2,List(foo3, foo4)))

}

【讨论】:

  • 好的,所以我需要做foreachPartition。我有关于 Elasticsearch 的其他问题,因为 Elasticsearch 只接受 RDD,所以我不能按分区做。我尝试使用 sc.MakeRDD 或 sc.parallize,都告诉我它不可序列化。请看我原来的问题。谢谢。
  • 您的批量输入是什么,以及您要对按键分组的值列表执行什么操作?
  • 我已经用更多细节更新了原始问题。谢谢!
猜你喜欢
  • 1970-01-01
  • 2020-12-11
  • 2018-03-30
  • 2016-05-20
  • 2017-04-11
  • 2023-03-30
  • 2021-06-17
  • 2022-01-01
  • 2018-08-29
相关资源
最近更新 更多