This is an automated email from the ASF dual-hosted git repository.

pjfanning pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko-connectors.git


The following commit(s) were added to refs/heads/main by this push:
     new 12b7e8267 Elasticsearch: support Elasticsearch 8 and 9 (#1934)
12b7e8267 is described below

commit 12b7e82679f1f5c9c87a8be878d0ce4e7240fd41
Author: PJ Fanning <[email protected]>
AuthorDate: Wed Sep 30 14:11:54 2026 +0100

    Elasticsearch: support Elasticsearch 8 and 9 (#1934)
    
    * Elasticsearch: support Elasticsearch 8 and 9
    
    Motivation:
    The connector only knew about `ApiVersion.V5` and `V7`, so users on
    Elasticsearch 8 or 9 had to declare `V7` and hope, and neither server
    line was covered by any test.
    
    Modification:
    Added `V8` and `V9` to `ApiVersion`, with matching
    `ElasticsearchParams.V8` / `V9` factories that take an index name only.
    Both dispatch to the existing V7 REST implementation: the bulk request
    already omits `_type`, the search endpoint is already
    `/<index>/_search`, and hit parsing reads only `_id`, `_source` and
    `_version`, all of which survive the removal of mapping types.
    
    Added `elasticsearch8` (8.19.21, port 9205) and `elasticsearch9`
    (9.5.3, port 9206) compose services with security disabled and a 512m
    heap, started them from the elasticsearch CI job, taught both the Scala
    and Java `constructElasticsearchParams` helpers about the new versions,
    and ran the existing `ElasticsearchConnectorBehaviour` suite against
    both servers. `createStrictMapping` in that suite now declares mappings
    without the `_doc` type and without `include_type_name` for V8 and V9,
    since Elasticsearch 8 removed both, and the two strict-mapping error
    assertions now match on the stable part of the reason rather than the
    trailing `within [_doc]`, which is server-internal.
    
    Result:
    Elasticsearch 8 and 9 users can name their server version, and both
    lines are covered by the full connector behaviour suite in Scala and by
    the parameterized test in Java. No existing signature changed, so
    binary compatibility is preserved.
    
    Tests:
    - `sbt "elasticsearch/Test/compile"` - success
    - `sbt "++3.3.8 elasticsearch/Test/compile"` - success
    - `sbt "elasticsearch/mimaReportBinaryIssues"` - success, no issues
    - `sbt "elasticsearch/scalafmtAll" "elasticsearch/javafmtAll"` - success
    - Integration run relies on the `elasticsearch` CI job; no local Docker.
    
    References:
    None - Elasticsearch 8 and 9 server support
    
    * Elasticsearch: stop the Scala and Java version-type tests sharing an index
    
    Motivation:
    The first CI run for Elasticsearch 8 and 9 support failed
    `ElasticsearchParameterizedTest.testUsingVersionType` on both new
    servers with `Unrecognized field "price"`. The Scala
    `ElasticsearchConnectorBehaviour` and the Java parameterized test both
    write to an index named `book-test-version-type` on the same server, and
    both write document id `1` at external version 5, so the second write
    loses to a version conflict and the reader sees the other suite's
    document. The Scala `Book` carries a `price` field the Java `Book`
    cannot deserialize. The race predates this change; running the behaviour
    suite against two more servers lost it.
    
    Modification:
    Renamed the Scala index to `book-test-version-type-scala`, matching the
    `version-test-scala` and `custom-search-params-test-scala` naming the
    same suite already uses for this reason.
    
    Result:
    The two suites no longer share an index, so neither can observe the
    other's documents.
    
    Tests:
    - `sbt "elasticsearch/Test/compile"` - success
    - `sbt "elasticsearch/scalafmtAll"` - success
    - Integration run relies on the `elasticsearch` CI job.
    
    References:
    Refs #1934
    
    * Elasticsearch: delete test indices by name instead of _all
    
    Motivation:
    Every test suite wipes its server between runs with `DELETE /_all`, and
    none of them check the response. Elasticsearch 8 changed
    `action.destructive_requires_name` to default to true, which refuses
    wildcard and `_all` destructive operations, so from Elasticsearch 8 on
    the cleanup silently does nothing. Documents then survive into the next
    suite, where they are read back by a suite that cannot deserialize them:
    the Java `Book` has only a `title`, while the Scala one also carries a
    `price`.
    
    Modification:
    Replaced the `_all` deletions with a shared helper that lists the
    indices via `_cat/indices` and deletes them by explicit name, which
    every supported server accepts. Both the list and the delete response
    are checked, so a cleanup that stops working fails the suite instead of
    passing quietly. Added it as `deleteAllIndices` in
    `ElasticsearchSpecUtils` for the Scala suites and rewrote
    `ElasticsearchTestBase.cleanIndex` for the Java ones. Also renamed the
    Scala `sink3-0` index to `sink3-0-scala`, the last index name shared
    between a Java and a Scala suite on the same server.
    
    Result:
    Cleanup works on every server the tests run against, and a future
    regression in it surfaces as a failure rather than as a confusing
    deserialization error in an unrelated suite.
    
    Tests:
    - `sbt "elasticsearch/Test/compile"` - success, no warnings
    - `sbt "++3.3.8 elasticsearch/Test/compile"` - success
    - `sbt "elasticsearch/scalafmtAll" "elasticsearch/javafmtAll"` - success
    - Integration run relies on the `elasticsearch` CI job; no local Docker.
    
    References:
    None - test isolation cleanup
    
    ---------
    
    Co-authored-by: Claude Opus 5 (1M context) <[email protected]>
---
 .github/workflows/check-build-test.yml             |  2 +-
 docker-compose.yml                                 | 16 +++++++
 docs/src/main/paradox/elasticsearch.md             |  8 ++--
 .../connectors/elasticsearch/ApiVersion.java       |  4 +-
 .../elasticsearch/ElasticsearchParams.scala        | 10 +++++
 .../impl/ElasticsearchSimpleFlowStage.scala        |  2 +-
 .../impl/ElasticsearchSourceStage.scala            |  5 ++-
 .../javadsl/ElasticsearchParameterizedTest.java    |  6 ++-
 .../java/docs/javadsl/ElasticsearchTestBase.java   | 49 +++++++++++++++++++++-
 .../scaladsl/ElasticsearchConnectorBehaviour.scala | 47 ++++++++++++---------
 .../scala/docs/scaladsl/ElasticsearchSpec.scala    | 25 +++++++----
 .../docs/scaladsl/ElasticsearchSpecUtils.scala     | 44 ++++++++++++++++++-
 .../scala/docs/scaladsl/ElasticsearchV5Spec.scala  |  8 +---
 .../scala/docs/scaladsl/ElasticsearchV7Spec.scala  |  8 +---
 .../test/scala/docs/scaladsl/OpensearchSpec.scala  |  8 +---
 15 files changed, 183 insertions(+), 59 deletions(-)

diff --git a/.github/workflows/check-build-test.yml 
b/.github/workflows/check-build-test.yml
index 5f58910a2..6c1144885 100644
--- a/.github/workflows/check-build-test.yml
+++ b/.github/workflows/check-build-test.yml
@@ -109,7 +109,7 @@ jobs:
           - { connector: couchbase3,                   pre_cmd: 'docker 
compose up -d couchbase8_prep' }
           - { connector: csv }
           - { connector: dynamodb,                     pre_cmd: 'docker 
compose up -d dynamodb' }
-          - { connector: elasticsearch,                pre_cmd: 'docker 
compose up -d elasticsearch6 elasticsearch7 opensearch2 opensearch3' }
+          - { connector: elasticsearch,                pre_cmd: 'docker 
compose up -d elasticsearch6 elasticsearch7 elasticsearch8 elasticsearch9 
opensearch2 opensearch3' }
           - { connector: file }
           - { connector: ftp,                          pre_cmd: 
'./scripts/ftp-servers.sh' }
           - { connector: geode,                        pre_cmd: 'docker 
compose up -d geode' }
diff --git a/docker-compose.yml b/docker-compose.yml
index 694a36b56..375ce234d 100644
--- a/docker-compose.yml
+++ b/docker-compose.yml
@@ -44,6 +44,22 @@ services:
       - "9202:9200"
     environment:
       - "discovery.type=single-node"
+  elasticsearch8:
+    image: docker.elastic.co/elasticsearch/elasticsearch:8.19.21
+    ports:
+      - "9205:9200"
+    environment:
+      - "discovery.type=single-node"
+      - "xpack.security.enabled=false"
+      - "ES_JAVA_OPTS=-Xms512m -Xmx512m"
+  elasticsearch9:
+    image: docker.elastic.co/elasticsearch/elasticsearch:9.5.3
+    ports:
+      - "9206:9200"
+    environment:
+      - "discovery.type=single-node"
+      - "xpack.security.enabled=false"
+      - "ES_JAVA_OPTS=-Xms512m -Xmx512m"
   opensearch2:
     image: opensearchproject/opensearch:2.19.6
     ports:
diff --git a/docs/src/main/paradox/elasticsearch.md 
b/docs/src/main/paradox/elasticsearch.md
index c34f565a7..861c7757b 100644
--- a/docs/src/main/paradox/elasticsearch.md
+++ b/docs/src/main/paradox/elasticsearch.md
@@ -124,7 +124,7 @@ Java
 | bufferSize             | 10             | `ElasticsearchSource` retrieves 
messages from Elasticsearch by scroll scan. This buffer size is used as the 
scroll size. | 
 | includeDocumentVersion | false          | Tell Elasticsearch to return the 
documents `_version` property with the search results. See 
[Version](https://www.elastic.co/guide/en/elasticsearch/reference/current/search-request-body.html#request-body-search-version)
 and [Optimistic Concurrenct 
Control](https://www.elastic.co/guide/en/elasticsearch/guide/current/optimistic-concurrency-control.html)
 to know about this property. |
 | scrollDuration         | 5 min          | `ElasticsearchSource`  retrieves 
messages from Elasticsearch by scroll scan. This parameter is used as a scroll 
value. See [Time 
units](https://www.elastic.co/guide/en/elasticsearch/reference/current/common-options.html#time-units)
 for supported units.                |
-| apiVersion             | V7             | Currently supports `V5` and `V7` 
(see below) |
+| apiVersion             | V7             | Currently supports `V5`, `V7`, 
`V8` and `V9` (see below) |
 
 ### Sink and flow configuration
 
@@ -142,7 +142,7 @@ Java
 | bufferSize          | 10         | Flow and Sink batch messages to bulk 
requests when back-pressure applies.                              |
 | versionType         | None       | If set, `ElasticsearchSink` uses the 
chosen versionType to index documents. See [Version 
types](https://www.elastic.co/guide/en/elasticsearch/reference/current/docs-index_.html#_version_types)
 for accepted settings. |
 | retryLogic          | No retries | See below |
-| apiVersion          | V7         | Currently supports `V5` and `V7` (see 
below) |
+| apiVersion          | V7         | Currently supports `V5`, `V7`, `V8` and 
`V9` (see below) |
 | allowExplicitIndex  | True       | When set to False, the index name will be 
included in the URL instead of on each document (see below) | 
 
 #### Retry logic
@@ -174,7 +174,9 @@ This will be used to:
 1. transform the bulk request into a format understood by the corresponding 
Elasticsearch server.
 2. determine whether to include the index type mapping in the API calls. See 
[removal of 
types](https://www.elastic.co/guide/en/elasticsearch/reference/current/removal-of-types.html)
 
-Currently 
[`V5`](https://www.elastic.co/guide/en/elasticsearch/reference/5.6/docs-bulk.html#docs-bulk)
 and 
[`V7`](https://www.elastic.co/guide/en/elasticsearch/reference/7.6/docs-bulk.html#docs-bulk)
 are supported specifically but this parameter does not need to match the 
server version exactly (for example, either `V5` or `V7` should work with 
Elasticsearch 6.x).
+Currently 
[`V5`](https://www.elastic.co/guide/en/elasticsearch/reference/5.6/docs-bulk.html#docs-bulk),
 
[`V7`](https://www.elastic.co/guide/en/elasticsearch/reference/7.6/docs-bulk.html#docs-bulk),
 `V8` and 
[`V9`](https://www.elastic.co/docs/api/doc/elasticsearch/operation/operation-bulk)
 are supported specifically but this parameter does not need to match the 
server version exactly (for example, either `V5` or `V7` should work with 
Elasticsearch 6.x).
+
+`V8` and `V9` share the type-less request shape that `V7` introduced, so they 
behave identically today; pick the one matching your server so the setting 
keeps meaning if the versions diverge. Elasticsearch 8 and 9 removed mapping 
types outright, so `ElasticsearchParams.V8` and `ElasticsearchParams.V9` take 
an index name only, like `ElasticsearchParams.V7`.
 
 ### Allow explicit index
 
diff --git 
a/elasticsearch/src/main/java/org/apache/pekko/stream/connectors/elasticsearch/ApiVersion.java
 
b/elasticsearch/src/main/java/org/apache/pekko/stream/connectors/elasticsearch/ApiVersion.java
index f5047b5f7..113bff93d 100644
--- 
a/elasticsearch/src/main/java/org/apache/pekko/stream/connectors/elasticsearch/ApiVersion.java
+++ 
b/elasticsearch/src/main/java/org/apache/pekko/stream/connectors/elasticsearch/ApiVersion.java
@@ -15,5 +15,7 @@ package org.apache.pekko.stream.connectors.elasticsearch;
 
 public enum ApiVersion implements 
org.apache.pekko.stream.connectors.elasticsearch.ApiVersionBase {
   V5,
-  V7
+  V7,
+  V8,
+  V9
 }
diff --git 
a/elasticsearch/src/main/scala/org/apache/pekko/stream/connectors/elasticsearch/ElasticsearchParams.scala
 
b/elasticsearch/src/main/scala/org/apache/pekko/stream/connectors/elasticsearch/ElasticsearchParams.scala
index 70fe1c0a1..6438d1fb1 100644
--- 
a/elasticsearch/src/main/scala/org/apache/pekko/stream/connectors/elasticsearch/ElasticsearchParams.scala
+++ 
b/elasticsearch/src/main/scala/org/apache/pekko/stream/connectors/elasticsearch/ElasticsearchParams.scala
@@ -25,6 +25,16 @@ object ElasticsearchParams {
     new ElasticsearchParams(indexName, None)
   }
 
+  /**
+   * Elasticsearch 8 dropped mapping types, so the params are shaped like the 
V7 ones.
+   */
+  def V8(indexName: String): ElasticsearchParams = V7(indexName)
+
+  /**
+   * Elasticsearch 9 dropped mapping types, so the params are shaped like the 
V7 ones.
+   */
+  def V9(indexName: String): ElasticsearchParams = V7(indexName)
+
   def V5(indexName: String, typeName: String): ElasticsearchParams = {
     require(indexName != null, "You must define an index name")
     require(typeName != null && typeName.trim.nonEmpty, "You must define a 
type name for ElasticSearch API version V5")
diff --git 
a/elasticsearch/src/main/scala/org/apache/pekko/stream/connectors/elasticsearch/impl/ElasticsearchSimpleFlowStage.scala
 
b/elasticsearch/src/main/scala/org/apache/pekko/stream/connectors/elasticsearch/impl/ElasticsearchSimpleFlowStage.scala
index 2aa2c2bb9..5cb97ed2d 100644
--- 
a/elasticsearch/src/main/scala/org/apache/pekko/stream/connectors/elasticsearch/impl/ElasticsearchSimpleFlowStage.scala
+++ 
b/elasticsearch/src/main/scala/org/apache/pekko/stream/connectors/elasticsearch/impl/ElasticsearchSimpleFlowStage.scala
@@ -53,7 +53,7 @@ private[elasticsearch] final class 
ElasticsearchSimpleFlowStage[T, C](
         settings.versionType,
         settings.allowExplicitIndex,
         writer)
-    case ApiVersion.V7 =>
+    case ApiVersion.V7 | ApiVersion.V8 | ApiVersion.V9 =>
       new RestBulkApiV7[T, C](elasticsearchParams.indexName, 
settings.versionType, settings.allowExplicitIndex, writer)
 
     case elasticsearch.OpensearchApiVersion.V1 =>
diff --git 
a/elasticsearch/src/main/scala/org/apache/pekko/stream/connectors/elasticsearch/impl/ElasticsearchSourceStage.scala
 
b/elasticsearch/src/main/scala/org/apache/pekko/stream/connectors/elasticsearch/impl/ElasticsearchSourceStage.scala
index 8c7f71a87..16e66c638 100644
--- 
a/elasticsearch/src/main/scala/org/apache/pekko/stream/connectors/elasticsearch/impl/ElasticsearchSourceStage.scala
+++ 
b/elasticsearch/src/main/scala/org/apache/pekko/stream/connectors/elasticsearch/impl/ElasticsearchSourceStage.scala
@@ -170,8 +170,9 @@ private[elasticsearch] final class 
ElasticsearchSourceLogic[T](
           val searchBody = 
ElasticsearchSourceStage.buildSearchBody(completeParams)
 
           val endpoint: String = settings.apiVersion match {
-            case ApiVersion.V5           => 
s"/${elasticsearchParams.indexName}/${elasticsearchParams.typeName.get}/_search"
-            case ApiVersion.V7           => 
s"/${elasticsearchParams.indexName}/_search"
+            case ApiVersion.V5                                 => 
s"/${elasticsearchParams.indexName}/${elasticsearchParams.typeName.get}/_search"
+            case ApiVersion.V7 | ApiVersion.V8 | ApiVersion.V9 =>
+              s"/${elasticsearchParams.indexName}/_search"
             case OpensearchApiVersion.V1 => 
s"/${elasticsearchParams.indexName}/_search"
             case other                   => throw new 
IllegalArgumentException(s"API version $other is not supported")
           }
diff --git 
a/elasticsearch/src/test/java/docs/javadsl/ElasticsearchParameterizedTest.java 
b/elasticsearch/src/test/java/docs/javadsl/ElasticsearchParameterizedTest.java
index 5df0b82dd..584962ff5 100644
--- 
a/elasticsearch/src/test/java/docs/javadsl/ElasticsearchParameterizedTest.java
+++ 
b/elasticsearch/src/test/java/docs/javadsl/ElasticsearchParameterizedTest.java
@@ -38,7 +38,11 @@ public class ElasticsearchParameterizedTest extends 
ElasticsearchTestBase {
   private ApiVersion apiVersion;
 
   public static Stream<Arguments> data() {
-    return Stream.of(Arguments.of(9201, ApiVersion.V5), Arguments.of(9202, 
ApiVersion.V7));
+    return Stream.of(
+        Arguments.of(9201, ApiVersion.V5),
+        Arguments.of(9202, ApiVersion.V7),
+        Arguments.of(9205, ApiVersion.V8),
+        Arguments.of(9206, ApiVersion.V9));
   }
 
   @AfterEach
diff --git 
a/elasticsearch/src/test/java/docs/javadsl/ElasticsearchTestBase.java 
b/elasticsearch/src/test/java/docs/javadsl/ElasticsearchTestBase.java
index 0a8c6edc9..06242e2b5 100644
--- a/elasticsearch/src/test/java/docs/javadsl/ElasticsearchTestBase.java
+++ b/elasticsearch/src/test/java/docs/javadsl/ElasticsearchTestBase.java
@@ -16,10 +16,12 @@ package docs.javadsl;
 import java.io.IOException;
 import java.util.ArrayList;
 import java.util.List;
+import java.util.stream.Collectors;
 import org.apache.pekko.actor.ActorSystem;
 import org.apache.pekko.http.javadsl.Http;
 import org.apache.pekko.http.javadsl.model.ContentTypes;
 import org.apache.pekko.http.javadsl.model.HttpRequest;
+import org.apache.pekko.http.javadsl.model.HttpResponse;
 import org.apache.pekko.stream.connectors.elasticsearch.ApiVersion;
 import org.apache.pekko.stream.connectors.elasticsearch.ApiVersionBase;
 import 
org.apache.pekko.stream.connectors.elasticsearch.ElasticsearchConnectionSettings;
@@ -79,9 +81,48 @@ public class ElasticsearchTestBase {
     flushAndRefresh("source");
   }
 
+  /**
+   * Deletes every index the tests created on the server.
+   *
+   * <p>{@code DELETE /_all} is refused from Elasticsearch 8 on, where {@code
+   * action.destructive_requires_name} defaults to true, so the indices are 
listed and deleted by
+   * explicit name instead. Both responses are checked, so a cleanup that 
silently stops working
+   * cannot leave documents behind for the next suite to trip over.
+   */
   protected static void cleanIndex() throws IOException {
-    HttpRequest request = 
HttpRequest.DELETE("%s/_all".formatted(connectionSettings.baseUrl()));
-    http.singleRequest(request).toCompletableFuture().join();
+    HttpRequest listRequest =
+        
HttpRequest.GET("%s/_cat/indices?h=index".formatted(connectionSettings.baseUrl()));
+    HttpResponse listResponse = 
http.singleRequest(listRequest).toCompletableFuture().join();
+    if (!listResponse.status().isSuccess()) {
+      throw new IOException("Failed to list indices: " + 
listResponse.status());
+    }
+
+    String body =
+        listResponse
+            .entity()
+            .toStrict(5000, system)
+            .toCompletableFuture()
+            .join()
+            .getData()
+            .utf8String();
+
+    List<String> indices =
+        body.lines()
+            .map(String::trim)
+            // leave the server's own bookkeeping indices alone
+            .filter(name -> !name.isEmpty() && !name.startsWith("."))
+            .collect(Collectors.toList());
+
+    if (!indices.isEmpty()) {
+      HttpRequest deleteRequest =
+          HttpRequest.DELETE(
+              "%s/%s".formatted(connectionSettings.baseUrl(), String.join(",", 
indices)));
+      HttpResponse deleteResponse = 
http.singleRequest(deleteRequest).toCompletableFuture().join();
+      if (!deleteResponse.status().isSuccess()) {
+        throw new IOException(
+            "Failed to delete indices " + indices + ": " + 
deleteResponse.status());
+      }
+    }
   }
 
   protected static void flushAndRefresh(String indexName) throws IOException {
@@ -157,6 +198,10 @@ public class ElasticsearchTestBase {
       return ElasticsearchParams.V5(indexName, typeName);
     } else if (apiVersion == ApiVersion.V7) {
       return ElasticsearchParams.V7(indexName);
+    } else if (apiVersion == ApiVersion.V8) {
+      return ElasticsearchParams.V8(indexName);
+    } else if (apiVersion == ApiVersion.V9) {
+      return ElasticsearchParams.V9(indexName);
     } else if (apiVersion == OpensearchApiVersion.V1) {
       return OpensearchParams.V1(indexName);
     } else {
diff --git 
a/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchConnectorBehaviour.scala
 
b/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchConnectorBehaviour.scala
index 65d368d2f..00fc59555 100644
--- 
a/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchConnectorBehaviour.scala
+++ 
b/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchConnectorBehaviour.scala
@@ -46,27 +46,32 @@ trait ElasticsearchConnectorBehaviour {
     import spray.json._
     import DefaultJsonProtocol._
 
+    // Elasticsearch 8 removed mapping types along with the 
`include_type_name` parameter,
+    // so from V8 on the mapping is declared without the enclosing `_doc` type.
+    val typelessMappings = apiVersion == ApiVersion.V8 || apiVersion == 
ApiVersion.V9
+
     def createStrictMapping(indexName: String): Unit = {
-      val uri = Uri(connectionSettings.baseUrl)
-        .withPath(Path(s"/$indexName"))
-        .withQuery(Uri.Query(Map("include_type_name" -> "true")))
+      val indexUri = 
Uri(connectionSettings.baseUrl).withPath(Path(s"/$indexName"))
+      val uri =
+        if (typelessMappings) indexUri
+        else indexUri.withQuery(Uri.Query(Map("include_type_name" -> "true")))
+
+      val mapping =
+        """{
+          |  "dynamic": "strict",
+          |  "properties": {
+          |    "title": { "type": "text"},
+          |    "price": { "type": "integer"}
+          |  }
+          |}""".stripMargin
+
+      val mappings = if (typelessMappings) mapping else s"""{ "_doc": $mapping 
}"""
 
       val request = HttpRequest(HttpMethods.PUT)
         .withUri(uri)
         .withEntity(
           ContentTypes.`application/json`,
-          s"""{
-            |  "mappings": {
-            |    "_doc": {
-            |      "dynamic": "strict",
-            |      "properties": {
-            |        "title": { "type": "text"},
-            |        "price": { "type": "integer"}
-            |      }
-            |    }
-            |  }
-            |}
-         """.stripMargin)
+          s"""{ "mappings": $mappings }""")
 
       http.singleRequest(request).futureValue
     }
@@ -210,8 +215,8 @@ trait ElasticsearchConnectorBehaviour {
         // Assert retired documents
         val failed = writeResults.filter(!_.success).head
         failed.message shouldBe WriteMessage.createIndexMessage("1", 
JsObject("subject" -> "Akka Concurrency".toJson))
-        failed.errorReason shouldBe Some(
-          "mapping set to strict, dynamic introduction of [subject] within 
[_doc] is not allowed")
+        failed.errorReason.get should include(
+          "mapping set to strict, dynamic introduction of [subject]")
 
         // Assert retried 5 times by looking duration
         assert(end - start > 5 * 100)
@@ -300,8 +305,8 @@ trait ElasticsearchConnectorBehaviour {
         WriteMessage
           .createIndexMessage("1", JsObject("subject" -> "Akka 
Concurrency".toJson))
           .withPassThrough(1)
-        failed.errorReason shouldBe Some(
-          "mapping set to strict, dynamic introduction of [subject] within 
[_doc] is not allowed")
+        failed.errorReason.get should include(
+          "mapping set to strict, dynamic introduction of [subject]")
 
         // Assert retried 5 times by looking duration
         assert(end - start > 5 * 100)
@@ -536,7 +541,9 @@ trait ElasticsearchConnectorBehaviour {
 
       "allow read and write using configured version type" in {
 
-        val indexName = "book-test-version-type"
+        // distinct from the Java suite's index of the same purpose: both 
suites share a server,
+        // and the Scala `Book` carries a `price` the Java `Book` cannot 
deserialize
+        val indexName = "book-test-version-type-scala"
         val typeName = "_doc"
 
         val book = Book("A sample title")
diff --git a/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchSpec.scala 
b/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchSpec.scala
index d8583e90e..0f32f0425 100644
--- a/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchSpec.scala
+++ b/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchSpec.scala
@@ -16,8 +16,6 @@ package docs.scaladsl
 import org.apache.pekko
 import pekko.actor.ActorSystem
 import pekko.http.scaladsl.{ Http, HttpExt }
-import pekko.http.scaladsl.model.Uri.Path
-import pekko.http.scaladsl.model.{ HttpMethods, HttpRequest, Uri }
 import pekko.stream.connectors.elasticsearch._
 import pekko.stream.connectors.testkit.scaladsl.LogCapturing
 import pekko.testkit.TestKit
@@ -43,15 +41,16 @@ class ElasticsearchSpec
     ElasticsearchConnectionSettings("http://localhost:9201";)
   val clientV7: ElasticsearchConnectionSettings =
     ElasticsearchConnectionSettings("http://localhost:9202";)
+  val clientV8: ElasticsearchConnectionSettings =
+    ElasticsearchConnectionSettings("http://localhost:9205";)
+  val clientV9: ElasticsearchConnectionSettings =
+    ElasticsearchConnectionSettings("http://localhost:9206";)
 
   override def afterAll(): Unit = {
-    val deleteRequestV5 = HttpRequest(HttpMethods.DELETE)
-      .withUri(Uri(clientV5.baseUrl).withPath(Path("/_all")))
-    http.singleRequest(deleteRequestV5).futureValue
-
-    val deleteRequestV7 = HttpRequest(HttpMethods.DELETE)
-      .withUri(Uri(clientV7.baseUrl).withPath(Path("/_all")))
-    http.singleRequest(deleteRequestV7).futureValue
+    deleteAllIndices(clientV5)
+    deleteAllIndices(clientV7)
+    deleteAllIndices(clientV8)
+    deleteAllIndices(clientV9)
 
     TestKit.shutdownActorSystem(system)
   }
@@ -64,4 +63,12 @@ class ElasticsearchSpec
     behave.like(elasticsearchConnector(ApiVersion.V7, clientV7))
   }
 
+  "Connector with ApiVersion 8 running against Elasticsearch v8" should {
+    behave.like(elasticsearchConnector(ApiVersion.V8, clientV8))
+  }
+
+  "Connector with ApiVersion 9 running against Elasticsearch v9" should {
+    behave.like(elasticsearchConnector(ApiVersion.V9, clientV9))
+  }
+
 }
diff --git 
a/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchSpecUtils.scala 
b/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchSpecUtils.scala
index 68379da86..99227c8e4 100644
--- a/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchSpecUtils.scala
+++ b/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchSpecUtils.scala
@@ -18,6 +18,7 @@ import pekko.actor.ActorSystem
 import pekko.http.scaladsl.HttpExt
 import pekko.http.scaladsl.model.Uri.Path
 import pekko.http.scaladsl.model.{ ContentTypes, HttpMethods, HttpRequest, Uri 
}
+import pekko.http.scaladsl.unmarshalling.Unmarshal
 import pekko.stream.connectors.elasticsearch.scaladsl.ElasticsearchSource
 import pekko.stream.connectors.elasticsearch.{
   ApiVersionBase,
@@ -31,7 +32,7 @@ import pekko.stream.scaladsl.Sink
 import org.scalatest.concurrent.ScalaFutures
 import org.scalatest.wordspec.AnyWordSpec
 
-import scala.concurrent.Future
+import scala.concurrent.{ ExecutionContext, Future }
 
 trait ElasticsearchSpecUtils { this: AnyWordSpec with ScalaFutures =>
   implicit def system: ActorSystem
@@ -57,6 +58,43 @@ trait ElasticsearchSpecUtils { this: AnyWordSpec with 
ScalaFutures =>
     http.singleRequest(request).futureValue
   }
 
+  /**
+   * Deletes every index the tests created on the server.
+   *
+   * `DELETE /_all` is refused from Elasticsearch 8 on, where 
`action.destructive_requires_name`
+   * defaults to true, so the indices are listed and deleted by explicit name 
instead. Both
+   * responses are checked, so a cleanup that silently stops working cannot 
leave documents behind
+   * for the next suite to trip over.
+   */
+  def deleteAllIndices(connectionSettings: ElasticsearchConnectionSettings): 
Unit = {
+    implicit val ec: ExecutionContext = system.dispatcher
+
+    val listRequest = HttpRequest(HttpMethods.GET)
+      .withUri(
+        Uri(connectionSettings.baseUrl)
+          .withPath(Path("/_cat/indices"))
+          .withQuery(Uri.Query("h" -> "index")))
+    val listResponse = http.singleRequest(listRequest).futureValue
+    require(listResponse.status.isSuccess(), s"Failed to list indices: 
${listResponse.status}")
+
+    val indices = Unmarshal(listResponse.entity)
+      .to[String]
+      .futureValue
+      .linesIterator
+      .map(_.trim)
+      // leave the server's own bookkeeping indices alone
+      .filter(name => name.nonEmpty && !name.startsWith("."))
+      .toList
+
+    if (indices.nonEmpty) {
+      val deleteRequest = HttpRequest(HttpMethods.DELETE)
+        .withUri(Uri(connectionSettings.baseUrl).withPath(Path("/" + 
indices.mkString(","))))
+      val deleteResponse = http.singleRequest(deleteRequest).futureValue
+      require(deleteResponse.status.isSuccess(),
+        s"Failed to delete indices ${indices.mkString(", ")}: 
${deleteResponse.status}")
+    }
+  }
+
   def flushAndRefresh(connectionSettings: ElasticsearchConnectionSettings, 
indexName: String): Unit = {
     val flushRequest = HttpRequest(HttpMethods.POST)
       
.withUri(Uri(connectionSettings.baseUrl).withPath(Path(s"/$indexName/_flush")))
@@ -98,6 +136,10 @@ trait ElasticsearchSpecUtils { this: AnyWordSpec with 
ScalaFutures =>
       ElasticsearchParams.V5(indexName, typeName)
     } else if (apiVersion == 
pekko.stream.connectors.elasticsearch.ApiVersion.V7) {
       ElasticsearchParams.V7(indexName)
+    } else if (apiVersion == 
pekko.stream.connectors.elasticsearch.ApiVersion.V8) {
+      ElasticsearchParams.V8(indexName)
+    } else if (apiVersion == 
pekko.stream.connectors.elasticsearch.ApiVersion.V9) {
+      ElasticsearchParams.V9(indexName)
     } else if (apiVersion == OpensearchApiVersion.V1) {
       OpensearchParams.V1(indexName)
     } else {
diff --git 
a/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchV5Spec.scala 
b/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchV5Spec.scala
index 65a79fb44..27e79cdd2 100644
--- a/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchV5Spec.scala
+++ b/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchV5Spec.scala
@@ -14,8 +14,6 @@
 package docs.scaladsl
 
 import org.apache.pekko
-import pekko.http.scaladsl.model.{ HttpMethods, HttpRequest, Uri }
-import pekko.http.scaladsl.model.Uri.Path
 import pekko.stream.connectors.elasticsearch.scaladsl.{ ElasticsearchFlow, 
ElasticsearchSink, ElasticsearchSource }
 import pekko.stream.connectors.elasticsearch.{
   ApiVersion,
@@ -47,9 +45,7 @@ class ElasticsearchV5Spec extends ElasticsearchSpecBase with 
ElasticsearchSpecUt
   }
 
   override def afterAll() = {
-    val deleteRequest = HttpRequest(HttpMethods.DELETE)
-      .withUri(Uri(connectionSettings.baseUrl).withPath(Path("/_all")))
-    http.singleRequest(deleteRequest).futureValue
+    deleteAllIndices(connectionSettings)
 
     TestKit.shutdownActorSystem(system)
   }
@@ -153,7 +149,7 @@ class ElasticsearchV5Spec extends ElasticsearchSpecBase 
with ElasticsearchSpecUt
     }
 
     "store properly formatted JSON from Strings" in {
-      val indexName = "sink3-0"
+      val indexName = "sink3-0-scala"
       // #string
       val write: Future[Seq[WriteResult[String, NotUsed]]] = Source(
         Seq(
diff --git 
a/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchV7Spec.scala 
b/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchV7Spec.scala
index cbb1a5ff8..562404d7f 100644
--- a/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchV7Spec.scala
+++ b/elasticsearch/src/test/scala/docs/scaladsl/ElasticsearchV7Spec.scala
@@ -14,8 +14,6 @@
 package docs.scaladsl
 
 import org.apache.pekko
-import pekko.http.scaladsl.model.Uri.Path
-import pekko.http.scaladsl.model.{ HttpMethods, HttpRequest, Uri }
 import pekko.stream.connectors.elasticsearch.scaladsl.{ ElasticsearchFlow, 
ElasticsearchSink, ElasticsearchSource }
 import pekko.stream.connectors.elasticsearch._
 import pekko.stream.scaladsl.{ Sink, Source }
@@ -39,9 +37,7 @@ class ElasticsearchV7Spec extends ElasticsearchSpecBase with 
ElasticsearchSpecUt
   }
 
   override def afterAll() = {
-    val deleteRequest = HttpRequest(HttpMethods.DELETE)
-      .withUri(Uri(connectionSettings.baseUrl).withPath(Path("/_all")))
-    http.singleRequest(deleteRequest).futureValue
+    deleteAllIndices(connectionSettings)
 
     TestKit.shutdownActorSystem(system)
   }
@@ -142,7 +138,7 @@ class ElasticsearchV7Spec extends ElasticsearchSpecBase 
with ElasticsearchSpecUt
     }
 
     "store properly formatted JSON from Strings" in {
-      val indexName = "sink3-0"
+      val indexName = "sink3-0-scala"
 
       val write: Future[Seq[WriteResult[String, NotUsed]]] = Source(
         Seq(
diff --git a/elasticsearch/src/test/scala/docs/scaladsl/OpensearchSpec.scala 
b/elasticsearch/src/test/scala/docs/scaladsl/OpensearchSpec.scala
index 050abed40..92766ca89 100644
--- a/elasticsearch/src/test/scala/docs/scaladsl/OpensearchSpec.scala
+++ b/elasticsearch/src/test/scala/docs/scaladsl/OpensearchSpec.scala
@@ -14,8 +14,6 @@
 package docs.scaladsl
 
 import org.apache.pekko
-import pekko.http.scaladsl.model.Uri.Path
-import pekko.http.scaladsl.model.{ HttpMethods, HttpRequest, Uri }
 import pekko.stream.connectors.elasticsearch.{
   ElasticsearchConnectionSettings,
   OpensearchApiVersion,
@@ -46,9 +44,7 @@ abstract class OpensearchSpec(baseUrl: String) extends 
ElasticsearchSpecBase wit
   }
 
   override def afterAll() = {
-    val deleteRequest = HttpRequest(HttpMethods.DELETE)
-      .withUri(Uri(connectionSettings.baseUrl).withPath(Path("/_all")))
-    http.singleRequest(deleteRequest).futureValue
+    deleteAllIndices(connectionSettings)
 
     TestKit.shutdownActorSystem(system)
   }
@@ -156,7 +152,7 @@ abstract class OpensearchSpec(baseUrl: String) extends 
ElasticsearchSpecBase wit
     }
 
     "store properly formatted JSON from Strings" in {
-      val indexName = "sink3-0"
+      val indexName = "sink3-0-scala"
 
       // #string
       val write: Future[Seq[WriteResult[String, NotUsed]]] = Source(


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to