【发布时间】:2018-06-02 05:26:43
【问题描述】:
我正在尝试从 Spark 在 Elasticsearch 中编写对象集合。我必须满足两个要求:
- 文档已经用 JSON 序列化,应该按原样编写
- 应提供 Elasticsearch 文档
_id
这是我目前尝试过的。
saveJsonToEs()
我尝试像这样使用saveJsonToEs()(序列化文档包含字段_id 和所需的Elasticsearch ID):
val rdd: RDD[String] = job.map{ r => r.toJson() }
val cfg = Map(
("es.resource", "myindex/mytype"),
("es.mapping.id", "_id"),
("es.mapping.exclude", "_id")
)
EsSpark.saveJsonToEs(rdd, cfg)
但是elasticsearch-hadoop 库给出了这个例外:
Caused by: org.elasticsearch.hadoop.EsHadoopIllegalArgumentException: When writing data as JSON, the field exclusion feature is ignored. This is most likely not what the user intended. Bailing out...
at org.elasticsearch.hadoop.util.Assert.isTrue(Assert.java:60)
at org.elasticsearch.hadoop.rest.InitializationUtils.validateSettings(InitializationUtils.java:253)
如果我删除es.mapping.exclude 但保留es.mapping.id 并发送带有_id 的JSON(如{"_id":"blah",...})
val cfg = Map(
("es.resource", "myindex/mytype"),
("es.mapping.id", "_id")
)
EsSpark.saveJsonToEs(rdd, cfg)
我收到此错误:
Exception in thread "main" org.apache.spark.SparkException: Job aborted due to stage failure: Task 15 in stage 84.0 failed 4 times, most recent failure: Lost task 15.3 in stage 84.0 (TID 628, 172.31.35.69, executor 1): org.apache.spark.util.TaskCompletionListenerException: Found unrecoverable error [172.31.30.184:9200] returned Bad Request(400) - Field [_id] is a metadata field and cannot be added inside a document. Use the index API request parameters.; Bailing out..
at org.apache.spark.TaskContextImpl.markTaskCompleted(TaskContextImpl.scala:105)
at org.apache.spark.scheduler.Task.run(Task.scala:112)
...
当我尝试将此 id 作为不同的字段发送时(例如 {"superID":"blah",...":
val cfg = Map(
("es.resource", "myindex/mytype"),
("es.mapping.id", "superID")
)
EsSpark.saveJsonToEs(rdd, cfg)
提取字段失败:
17/12/20 15:15:38 WARN TaskSetManager: Lost task 8.0 in stage 84.0 (TID 586, 172.31.33.56, executor 0): org.elasticsearch.hadoop.EsHadoopIllegalArgumentException: [JsonExtractor for field [superId]] cannot extract value from entity [class java.lang.String] | instance [{...,"superID":"7f48c8ee6a8a"}]
at org.elasticsearch.hadoop.serialization.bulk.AbstractBulkFactory$FieldWriter.write(AbstractBulkFactory.java:106)
at org.elasticsearch.hadoop.serialization.bulk.TemplatedBulk.writeTemplate(TemplatedBulk.java:80)
at org.elasticsearch.hadoop.serialization.bulk.TemplatedBulk.write(TemplatedBulk.java:56)
at org.elasticsearch.hadoop.rest.RestRepository.writeToIndex(RestRepository.java:161)
at org.elasticsearch.spark.rdd.EsRDDWriter.write(EsRDDWriter.scala:67)
at org.elasticsearch.spark.rdd.EsSpark$$anonfun$doSaveToEs$1.apply(EsSpark.scala:107)
at org.elasticsearch.spark.rdd.EsSpark$$anonfun$doSaveToEs$1.apply(EsSpark.scala:107)
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87)
当我从配置中删除 es.mapping.id 和 es.mapping.exclude 时,它可以工作,但文档 ID 是由 Elasticsearch 生成的(这违反了要求 2):
val rdd: RDD[String] = job.map{ r => r.toJson() }
val cfg = Map(
("es.resource", "myindex/mytype"),
)
EsSpark.saveJsonToEs(rdd, cfg)
saveToEsWithMeta()
还有另一个函数可以提供_id 和其他metadata 用于插入:saveToEsWithMeta() 允许解决要求 2 但因要求 1 而失败。
val rdd: RDD[(String, String)] = job.map{
r => r._id -> r.toJson()
}
val cfg = Map(
("es.resource", "myindex/mytype"),
)
EsSpark.saveToEsWithMeta(rdd, cfg)
事实上,Elasticsearch 甚至无法解析 elasticsearch-hadoop 发送的内容:
Caused by: org.apache.spark.util.TaskCompletionListenerException: Found unrecoverable error [<es_host>:9200] returned Bad Request(400) - failed to parse; Bailing out..
at org.apache.spark.TaskContextImpl.markTaskCompleted(TaskContextImpl.scala:105)
at org.apache.spark.scheduler.Task.run(Task.scala:112)
问题
是否可以将 Spark 中的 (documentID, serializedDocument) 集合写入 Elasticsearch(使用 elasticsearch-hadoop)?
附:我正在使用 Elasticsearch 5.6.3 和 Spark 2.1.1。
【问题讨论】:
-
您是否尝试在删除
es.mapping.exclude的同时保留es.mapping.id设置?老实说,我什至不明白你为什么需要排除 -
@AlexSavitsky 感谢您的注意,我更新了问题!但它仍然不起作用:(
-
我建议尝试这样做的原因是我最近有一个类似的用例,只有我的数据是具有明确定义的架构的 CSV。我的 ID 列是架构的一部分,它也在
es.mapping.id中指定,并且它被保存到 ES 中没有问题。我认为这对 JSON 数据同样适用,但显然情况并非如此。
标签: scala apache-spark elasticsearch elasticsearch-hadoop