【发布时间】: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 以动态获取值。
【问题讨论】:
-
查看下面的例子,你会更好理解。
-
ES 未针对许多小型插入进行优化。如果可以尝试暂时关闭 ES 中的索引,保存数据并再次打开索引。 elastic.co/guide/en/elasticsearch/reference/current/…
标签: scala apache-spark elasticsearch