【问题标题】:How to export all data from Elastic Search Index to file in JSON format with _id field specified?如何将 Elastic Search Index 中的所有数据以 JSON 格式导出到指定 _id 字段的文件?
【发布时间】:2019-07-09 07:17:44
【问题描述】:

我是 Spark 和 Scala 的新手。我正在尝试将 Elastic Search 中特定索引中的所有数据读取到 RDD 中,并使用这些数据写入 Mongo DB。

我正在将 Elastic 搜索数据加载到 esJsonRDD,当我尝试打印 RDD 内容时,它采用以下格式,

(1765770532{"FirstName":ABC,"LastName":"DEF",Zipcode":"36905","City":"PortAdam","StateCode":"AR"})

预期格式,

{_id:"1765770532","FirstName":ABC,"LastName":"DEF",Zipcode":"36905","City":"PortAdam","StateCode":"AR"}

如何实现弹性搜索的输出以这种方式格式化?

任何帮助将不胜感激。

elasticsearch检索到的数据格式如下,

(1765770532{"FirstName":ABC,"LastName":"DEF",Zipcode":"36905","City":"PortAdam","StateCode":"AR"})

预期格式是,

{_id:"1765770532","FirstName":ABC,"LastName":"DEF",Zipcode":"36905","City":"PortAdam","StateCode":"AR"}

    object readFromES {

    def main(args: Array[String]) {

        val conf = new SparkConf().setAppName("readFromES")
        .set("es.nodes", Config.ES_NODES)
        .set("es.nodes.wan.only", Config.ES_NODES_WAN_ONLY)
        .set("es.net.http.auth.user", Config.ES_NET_HTTP_AUTH_USER)
        .set("es.net.http.auth.pass", Config.ES_NET_HTTP_AUTH_PASS)
        .set("es.net.ssl", Config.ES_NET_SSL)
        .set("es.output.json","true")

        val sc = new SparkContext(conf)
        val RDD =  EsSpark.esJsonRDD(sc, "userdata/user")
        //RDD.coalesce(1).saveAsTextFile(args(0))
        RDD.take(5).foreach(println)
        }
       }

我希望将 RDD 输出写入以下 JSON 格式的文件(每个文档一行),

{_id:"1765770532","FirstName":ABC,"LastName":"DEF",Zipcode":"36905","City":"PortAdam","StateCode":"AR"}
{_id:"1765770533","FirstName":DEF,"LastName":"DEF",Zipcode":"35525","City":"PortWinchestor","StateCode":"AI"}

【问题讨论】:

  • 你用的是什么版本的spark?
  • Spark 版本 2.2.1

标签: json scala apache-spark elasticsearch


【解决方案1】:

"_id" 是元数据的一部分,要访问它,您应该将.config("es.read.metadata", true) 添加到配置中。

那么你可以通过两种方式访问​​它,你可以使用

val RDD =  EsSpark.esJsonRDD(sc, "userdata/user") 

并在json中手动添加_id字段

或者更简单的方法是读取为数据框

val df = spark.read
  .format("org.elasticsearch.spark.sql")
  .load("userdata/user")
  .withColumn("_id", $"_metadata".getItem("_id"))
  .drop("_metadata")

//写成json到文件中

df.write.json("output folder ")

这里的 spark 是创建为

的 spark 会话
val spark = SparkSession.builder().master("local[*]").appName("Test")
  .config("spark.es.nodes","host")
  .config("spark.es.port","ports")
  .config("spark.es.nodes.wan.only","true")
  .config("es.read.metadata", true) //for enabling metadata
  .getOrCreate()

希望对你有帮助

【讨论】:

  • 感谢您的回复。我需要在以下方面进一步澄清,1.说,我想在从索引读取时根据ES查询过滤数据,我该如何以数据框的方式做到这一点? DF 是否支持这样的东西,val RDD = EsSpark.esJsonRDD(sc, "UserData/user",myQuery) 2. 说,我想将此数据帧转换为 JSONRDD,然后使用它将其推送到 Mongo DB,如何我能做到吗?提前致谢
  • 对于 1,您可以添加查询 .option("query", myQuery) 并在配置中使用 .option("pushdown", "true")。 2 您可以使用mongo-spark-connector直接将数据框添加到mongo集合中
  • 如果以上答案有帮助,请标记为答案并关闭它,您已经问了几个其他问题,将这些问题移到另一个问题
  • 我当然会这样做。我面临以下错误,在尝试使用您上面提到的方法时,值 $ 不是 StringContext.withColumn("_id", $"_metadata".getItem("_id")) 的成员
  • 你需要导入import spark.implicits._或者你可以使用col("_metadata")代替$"_metadata"
猜你喜欢
  • 2016-01-01
  • 2018-03-29
  • 2017-04-08
  • 1970-01-01
  • 1970-01-01
  • 2020-04-27
  • 1970-01-01
  • 1970-01-01
  • 2018-01-14
相关资源
最近更新 更多