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]