【问题标题】:NullPointerException while using elastic search 5.5 Bulk Ingest API with Spark Streaming将弹性搜索 5.5 Bulk Ingest API 与 Spark Streaming 一起使用时出现 NullPointerException
【发布时间】:2018-07-24 16:24:31
【问题描述】:

获取 NullPointerException:

java.lang.NullPointerException
    at org.elasticsearch.action.bulk.BulkRequest.validate(BulkRequest.java:604)
    at org.elasticsearch.action.TransportActionNodeProxy.execute(TransportActionNodeProxy.java:46)
    at org.elasticsearch.client.transport.TransportProxyClient.lambda$execute$0(TransportProxyClient.java:59)
    at org.elasticsearch.client.transport.TransportClientNodesService.execute(TransportClientNodesService.java:250)
    at org.elasticsearch.client.transport.TransportProxyClient.execute(TransportProxyClient.java:59)
    at org.elasticsearch.client.transport.TransportClient.doExecute(TransportClient.java:363)
    at org.elasticsearch.client.support.AbstractClient.execute(AbstractClient.java:408)
    at org.elasticsearch.action.ActionRequestBuilder.execute(ActionRequestBuilder.java:80)
    at org.elasticsearch.action.ActionRequestBuilder.execute(ActionRequestBuilder.java:54)

我有一个场景,多个并发任务在 Spark Streaming Application 的 4 个执行程序中运行,每个执行程序都从 Kafka 读取数据,准备批量并摄取 ES 索引中的批量记录。我第一次收到这些奇怪的 NullPointerException 与其中一些记录,但它们在第二次运行中被成功处理。

谁能告诉我为什么会这样。

【问题讨论】:

  • 您是否考虑过仅使用 Kafka Connect(它是 Apache Kafka 的一部分)将数据从 Kafka 流式传输到 Elasticsearch?看看:speakerdeck.com/rmoff/…
  • 到目前为止,我们正在使用 Spark Kafka Streaming 作为我们当前工具堆栈的一部分。但这确实是您分享的内容,将看看它。谢谢!

标签: scala apache-spark elasticsearch apache-kafka


【解决方案1】:

这是我使用的代码 sn-p 第一行是来自我的 build.sbt 文件的依赖项

//lib dependency in build.sbt
"org.elasticsearch" %% "elasticsearch-spark-20" % "5.6.5"

//below is the connection variables required by Spark

val resources: String =
  s"${appConf.getString("es-index")}/${appConf.getString("es.type")}"
val esConfig: Map[String, String] = Map(
  "es.index.auto.create" -> s"${appConf.getString("es.index.auto.create")}",
  "es.nodes" -> s"${appConf.getString("es-nodes")}",
  "es.port" -> s"${appConf.getInt("es.port")}",
  "es.nodes.wan.only" -> s"${appConf.getString("es.nodes.wan.only")}",
  "es.net.ssl" -> s"${appConf.getString("es.net.ssl")}"
)

import org.elasticsearch.spark._
    val dstream: InputDStream[ConsumerRecord[String, String]] =
  KafkaUtils.createDirectStream[String, String](
    ssc,
    LocationStrategies.PreferConsistent,
    ConsumerStrategies.Subscribe[String, String](conn.topic,
                                                 conn.kafkaProps)
  )
dstream.foreachRDD(rdd =>
  rdd.map(_.value).saveJsonToEs(resources,esConfig))
ssc.checkpoint("/tmp/OACSpark")
ssc.start()
ssc.awaitTermination()

我使用类型安全配置从属性文件中读取配置。 我以 json 的形式向 kafka 发布数据,所以我使用了“saveJsonToEs()”api,您可以在 Elasticsearch 网站上的连接器文档中找到更多信息

【讨论】:

  • 谢谢@Yayati...但到目前为止我们还没有使用 Spark ES 连接器。请让我知道是否有办法通过我们目前使用的 ES JAVA API 来实现。
  • 嗨,肖比特。 ES-Spark 连接器在底层使用 ES Java 连接器。我相信你会没事的。
  • 嘿Yayati..我们正计划继续。谢谢!
【解决方案2】:

到目前为止,我有一个解决方法,可以一次将记录推送到 ES 索引,并删除了这个批量 API(批量 API 在后台也做同样的事情)。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-07-03
    • 2018-10-07
    • 1970-01-01
    • 2018-03-26
    • 1970-01-01
    • 1970-01-01
    • 2020-11-17
    • 1970-01-01
    相关资源
    最近更新 更多