【问题标题】:Elasticsearch + Spark: write json with custom document _idElasticsearch + Spark:使用自定义文档_id编写json
【发布时间】:2018-06-02 05:26:43
【问题描述】:

我正在尝试从 Spark 在 Elasticsearch 中编写对象集合。我必须满足两个要求:

  1. 文档已经用 JSON 序列化,应该按原样编写
  2. 应提供 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.ides.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


【解决方案1】:

最后我发现了问题:这是配置中的一个错字。

[JsonExtractor for field [superId]] cannot extract value from entity [class java.lang.String] | instance [{...,"superID":"7f48c8ee6a8a"}]

它正在寻找一个字段superID,但只有superID(注意大小写)。在问题中它也有点误导,因为在代码中它看起来像"es.mapping.id", "superID"(这是不正确的)。

实际解决方案如Levi Ramsey建议:

val json = """{"foo":"bar","superID":"deadbeef"}"""

val rdd = spark.makeRDD(Seq(json))
val cfg = Map(
  ("es.mapping.id", "superID"),
  ("es.resource", "myindex/mytype")
)
EsSpark.saveJsonToEs(rdd, cfg = cfg)

区别在于es.mapping.id不能是_id(如原帖所示,_id是元数据,Elasticsearch不接受)。

这自然意味着应该将新字段superID 添加到映射中(除非映射是动态的)。如果在索引中存储额外的字段是一种负担,还应该:

  • exclude它来自映射
  • 并禁用其索引

非常感谢Alex Savitsky 指出正确的方向。

【讨论】:

    【解决方案2】:

    您是否尝试过类似的方法:

    val rdd: RDD[String] = job.map{ r => r.toJson() }
    val cfg = Map(
      ("es.mapping.id", "_id")
    )
    rdd.saveJsonToEs("myindex/mytype", cfg)
    

    我已经针对 ES 1.7 进行了测试(使用 elasticsearch-hadoop(连接器版本 2.4.5))并且可以正常工作。

    【讨论】:

      【解决方案3】:

      可以通过将ES_INPUT_JSON 选项传递给cfg 来完成 参数映射并返回一个元组,其中包含文档 id 作为第一个元素,并在 JSON 中序列化的文档作为映射函数的第二个元素。

      我使用 "org.elasticsearch" %% "elasticsearch-spark-20" % "[6.0,7.0[" 对 Elasticsearch 6.4 进行了测试

      import org.elasticsearch.hadoop.cfg.ConfigurationOptions.{ES_INPUT_JSON, ES_NODES}
      import org.elasticsearch.spark._
      import org.elasticsearch.spark.sql._
      
      job
        .map{ r => (r._id, r.toJson()) }
        .saveToEsWithMeta(
          "myindex/mytype",
          Map(
            ES_NODES -> "https://localhost:9200",
            ES_INPUT_JSON -> true.toString
          )
        )
      

      【讨论】:

        【解决方案4】:

        我花了几天的时间把头撞在墙上,试图弄清楚为什么当我像这样使用字符串作为 ID 时 saveToEsWithMeta 不起作用:

        rdd.map(caseClassContainingJson =>
          (caseClassContainingJson._idWhichIsAString, caseClassContainingJson.jsonString)
        )
        .saveToEsWithMeta(s"$nationalShapeIndexName/$nationalShapeIndexType", Map(
          ES_INPUT_JSON -> true.toString
        ))
        

        这将引发与 JSON 解析相关的错误,从而使您误以为问题出在 JSON 上,但随后您记录每个 JSON 并查看它们都是有效的。

        事实证明,无论出于何种原因,ES_INPUT_JSON -&gt; true 使元组的左侧(即 ID)也被解析为 JSON!

        解决方案,JSON 将 ID 字符串化(将 ID 用额外的双引号括起来),以便将其解析为 JSON:

        rdd.map(caseClassContainingJson =>
          (
            Json.stringify(JsString(caseClassContainingJson._idWhichIsAString)), 
            caseClassContainingJson.jsonString
          )
        )
        .saveToEsWithMeta(s"$nationalShapeIndexName/$nationalShapeIndexType", Map(
          ES_INPUT_JSON -> true.toString
        ))
        

        【讨论】:

          【解决方案5】:
          1. 您可以使用saveToEs 来定义customer_id 而不必保存customer_id
          2. 注意rdd是RDD[Map]类型
          val rdd:RDD[Map[String, Any]]=...
          val cfg = Map(
            ("es.mapping.id", your_customer_id),
            ("es.mapping.exclude", your_customer_id)
          )
          EsSpark.saveToEs(rdd, your_es_index, cfg)
          

          【讨论】:

            猜你喜欢
            • 1970-01-01
            • 2019-12-06
            • 1970-01-01
            • 2021-07-22
            • 1970-01-01
            • 1970-01-01
            • 2015-07-29
            • 1970-01-01
            • 2017-06-22
            相关资源
            最近更新 更多