Akka-HTTP调用Elasticsearch 8.7.0 _bulk接口报错:找不到处理程序
问题描述
我使用Elasticsearch 8.7.0,已通过以下curl命令创建myindex索引:
curl --insecure --user "elastic:password" -X PUT "https://localhost:9200/myindex"
同时可通过curl手动向Elasticsearch 8.7.0添加JSON文档:
curl --insecure --user "elastic:password" -X POST "https://localhost:9200/myindex/_bulk" -H 'Content-Type: application/json' -d'{"some": "json"}'
但尝试用Scala 3结合Akka-HTTP通过以下代码编程调用时:
@main def main(args: String*): Unit = { val records = Seq( """{ "index":{ "_index" : "myindex", "_id" : "1" } } |{"title":"title1"}""".stripMargin, """{ "index":{ "_index" : "myindex", "_id" : "2" } } |{"title":"title2"}""".stripMargin, """{ "index":{ "_index" : "myindex", "_id" : "3" } } |{"title":"title3"}""".stripMargin, """{ "index":{ "_index" : "myindex", "_id" : "4" } } |{"title":"title4"}""".stripMargin) implicit val system: ActorSystem = ActorSystem("elastic") implicit val ec: ExecutionContextExecutor = system.dispatcher val trustfulSslContext: SSLContext = { object NoCheckX509TrustManager extends X509TrustManager { override def checkClientTrusted(chain: Array[X509Certificate], authType: String): Unit = () override def checkServerTrusted(chain: Array[X509Certificate], authType: String): Unit = () override def getAcceptedIssuers: Array[X509Certificate] = Array[X509Certificate]() } val context: SSLContext = SSLContext.getInstance("TLS") context.init(Array[KeyManager](), Array(NoCheckX509TrustManager), new SecureRandom()) context } Source .fromIterator[String]( () => records.iterator ) .grouped(2) .map(_.mkString("\n")) .map { payload => val entity = HttpEntity(ContentTypes.`application/json`, payload) val request = HttpRequest( HttpMethods.POST, uri = Uri("https://localhost:9200/myindex/_bulk"), headers = Seq(headers.Authorization(BasicHttpCredentials("elastic", "password"))), entity = entity ) println(s"request: $request") request } .via(Http().outgoingConnectionHttps(host = "localhost", port = 9200, connectionContext = ConnectionContext.httpsClient(trustfulSslContext))) .map(response => { println(s"response: ${response.entity.toStrict(1 second).map(_.data.utf8String)}") response }) .toMat(Sink.ignore)(Keep.none) .run() }
得到如下输出:
request: HttpRequest(HttpMethod(POST),https://localhost:9200/myindex/_bulk,List(Authorization),HttpEntity.Strict(application/json,137 bytes total),HttpProtocol(HTTP/1.1)) response: FulfilledFuture({"error":"no handler found for uri [https://localhost:9200/myindex/_bulk] and method [POST]"}) request: HttpRequest(HttpMethod(POST),https://localhost:9200/myindex/_bulk,List(Authorization),HttpEntity.Strict(application/json,137 bytes total),HttpProtocol(HTTP/1.1)) response: FulfilledFuture({"error":"no handler found for uri [https://localhost:9200/myindex/_bulk] and method [POST]"})
返回错误信息:{"error":"no handler found for uri [https://localhost:9200/myindex/_bulk] and method [POST]"},请问我遗漏了什么?
解决方案
问题出在HttpRequest的URI设置上:当使用Akka-HTTP的outgoingConnectionHttps时,你已经指定了连接的host和port(localhost:9200),此时HttpRequest的uri应该只填写路径部分,而不是完整的URL。
如果你传入完整的https://localhost:9200/myindex/_bulk,Akka-HTTP会把整个URL当作请求路径发送给Elasticsearch,导致Elasticsearch找不到对应的处理程序(它期望的路径是/myindex/_bulk,而不是完整URL)。
修改步骤:
- 将HttpRequest中的uri从完整URL改为仅路径:
uri = Uri("/myindex/_bulk")
- 额外补充:Elasticsearch的bulk API要求每一行数据(包括最后一行)都必须以换行符结尾,否则可能解析失败。你当前的
mkString("\n")会把两个record用换行分隔,但每个record内部已经有换行,建议在每个payload末尾再加一个换行:
.map(_.mkString("\n") + "\n")
- 最后记得在程序结束时关闭ActorSystem,避免资源泄漏:
// 在run()之后添加 sys.addShutdownHook { system.terminate() }
修改后的核心代码片段:
.map { payload => val finalPayload = payload + "\n" // 确保末尾有换行 val entity = HttpEntity(ContentTypes.`application/json`, finalPayload) val request = HttpRequest( HttpMethods.POST, uri = Uri("/myindex/_bulk"), // 仅路径部分 headers = Seq(headers.Authorization(BasicHttpCredentials("elastic", "password"))), entity = entity ) println(s"request: $request") request }
内容的提问来源于stack exchange,提问作者david
相关产品推荐
相关产品推荐

