From e1c3a5a3cc6a18d5ca7a8128a6b8ee0c785d2744 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Wed, 28 Jun 2023 10:11:56 +0300 Subject: [PATCH 1/5] [improve][test] Disable geoip download and enable logging for Elastic Testcontainers --- .../elasticsearch/ElasticSearchTestBase.java | 28 +++++++++++------ .../io/sinks/ElasticSearch7SinkTester.java | 4 +-- .../io/sinks/ElasticSearch8SinkTester.java | 3 +- .../io/sinks/ElasticSearchSinkTester.java | 30 +++++++++++++------ .../io/sinks/OpenSearchSinkTester.java | 9 ++---- 5 files changed, 46 insertions(+), 28 deletions(-) diff --git a/pulsar-io/elastic-search/src/test/java/org/apache/pulsar/io/elasticsearch/ElasticSearchTestBase.java b/pulsar-io/elastic-search/src/test/java/org/apache/pulsar/io/elasticsearch/ElasticSearchTestBase.java index 4c6fd020fa338..9e50a763cf061 100644 --- a/pulsar-io/elastic-search/src/test/java/org/apache/pulsar/io/elasticsearch/ElasticSearchTestBase.java +++ b/pulsar-io/elastic-search/src/test/java/org/apache/pulsar/io/elasticsearch/ElasticSearchTestBase.java @@ -18,10 +18,6 @@ */ package org.apache.pulsar.io.elasticsearch; -import java.io.IOException; -import java.util.Map; -import java.util.Optional; - import co.elastic.clients.elasticsearch.ElasticsearchClient; import co.elastic.clients.elasticsearch.security.CreateApiKeyRequest; import co.elastic.clients.elasticsearch.security.CreateApiKeyResponse; @@ -29,6 +25,10 @@ import co.elastic.clients.elasticsearch.security.GetTokenResponse; import co.elastic.clients.elasticsearch.security.get_token.AccessTokenGrantType; import com.fasterxml.jackson.databind.ObjectMapper; +import java.io.IOException; +import java.util.Map; +import java.util.Optional; +import lombok.extern.slf4j.Slf4j; import org.apache.pulsar.io.elasticsearch.client.elastic.ElasticSearchJavaRestClient; import org.apache.pulsar.io.elasticsearch.client.opensearch.OpenSearchHighLevelRestClient; import org.opensearch.client.Request; @@ -36,6 +36,7 @@ import org.testcontainers.elasticsearch.ElasticsearchContainer; import org.testcontainers.utility.DockerImageName; +@Slf4j public abstract class ElasticSearchTestBase { public static final String ELASTICSEARCH_8 = Optional.ofNullable(System.getenv("ELASTICSEARCH_IMAGE_V8")) @@ -54,17 +55,26 @@ public ElasticSearchTestBase(String elasticImageName) { } protected ElasticsearchContainer createElasticsearchContainer() { + ElasticsearchContainer elasticsearchContainer; if (elasticImageName.equals(OPENSEARCH)) { DockerImageName dockerImageName = DockerImageName.parse(OPENSEARCH).asCompatibleSubstituteFor("docker.elastic.co/elasticsearch/elasticsearch"); - return new ElasticsearchContainer(dockerImageName) + elasticsearchContainer = new ElasticsearchContainer(dockerImageName) .withEnv("OPENSEARCH_JAVA_OPTS", "-Xms128m -Xmx256m") .withEnv("bootstrap.memory_lock", "true") .withEnv("plugins.security.disabled", "true"); + } else { + elasticsearchContainer = new ElasticsearchContainer(elasticImageName) + .withEnv("ES_JAVA_OPTS", "-Xms128m -Xmx256m") + .withEnv("xpack.security.enabled", "false") + .withEnv("xpack.security.http.ssl.enabled", "false"); } - return new ElasticsearchContainer(elasticImageName) - .withEnv("ES_JAVA_OPTS", "-Xms128m -Xmx256m") - .withEnv("xpack.security.enabled", "false") - .withEnv("xpack.security.http.ssl.enabled", "false"); + configureElasticContainer(elasticsearchContainer); + return elasticsearchContainer; + } + + protected void configureElasticContainer(ElasticsearchContainer elasticContainer) { + elasticContainer.withEnv("ingest.geoip.downloader.enabled", "false") + .withLogConsumer(o -> log.info("elastic> {}", o.getUtf8String())); } protected ElasticSearchConfig.CompatibilityMode getCompatibilityMode() { diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearch7SinkTester.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearch7SinkTester.java index 65b38c677bfc5..d99fcad252706 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearch7SinkTester.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearch7SinkTester.java @@ -19,7 +19,6 @@ package org.apache.pulsar.tests.integration.io.sinks; import java.util.Optional; -import org.apache.pulsar.tests.integration.topologies.PulsarCluster; import org.testcontainers.elasticsearch.ElasticsearchContainer; public class ElasticSearch7SinkTester extends ElasticSearchSinkTester { @@ -32,8 +31,9 @@ public ElasticSearch7SinkTester(boolean schemaEnable) { super(schemaEnable); } + @Override - protected ElasticsearchContainer createSinkService(PulsarCluster cluster) { + protected ElasticsearchContainer createElasticContainer() { return new ElasticsearchContainer(ELASTICSEARCH_7) .withEnv("ES_JAVA_OPTS", "-Xms128m -Xmx256m"); } diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearch8SinkTester.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearch8SinkTester.java index bb52c4ff03fea..6ea7d3c1246f8 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearch8SinkTester.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearch8SinkTester.java @@ -19,7 +19,6 @@ package org.apache.pulsar.tests.integration.io.sinks; import java.util.Optional; -import org.apache.pulsar.tests.integration.topologies.PulsarCluster; import org.testcontainers.elasticsearch.ElasticsearchContainer; public class ElasticSearch8SinkTester extends ElasticSearchSinkTester { @@ -33,7 +32,7 @@ public ElasticSearch8SinkTester(boolean schemaEnable) { } @Override - protected ElasticsearchContainer createSinkService(PulsarCluster cluster) { + protected ElasticsearchContainer createElasticContainer() { return new ElasticsearchContainer(ELASTICSEARCH_8) .withEnv("ES_JAVA_OPTS", "-Xms128m -Xmx256m") .withEnv("xpack.security.enabled", "false") diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearchSinkTester.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearchSinkTester.java index 546dd1b9113ab..c7d8c84f684a0 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearchSinkTester.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearchSinkTester.java @@ -19,15 +19,6 @@ package org.apache.pulsar.tests.integration.io.sinks; import static org.testng.Assert.assertTrue; - -import java.util.Arrays; -import java.util.HashMap; -import java.util.HashSet; -import java.util.LinkedHashMap; -import java.util.List; -import java.util.Map; -import java.util.Set; - import co.elastic.clients.elasticsearch.ElasticsearchClient; import co.elastic.clients.elasticsearch.core.SearchRequest; import co.elastic.clients.elasticsearch.core.SearchResponse; @@ -35,6 +26,13 @@ import co.elastic.clients.transport.ElasticsearchTransport; import co.elastic.clients.transport.rest_client.RestClientTransport; import com.google.common.collect.ImmutableMap; +import java.util.Arrays; +import java.util.HashMap; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; import lombok.AllArgsConstructor; import lombok.Cleanup; import lombok.Data; @@ -46,6 +44,7 @@ import org.apache.pulsar.common.schema.KeyValue; import org.apache.pulsar.common.schema.KeyValueEncodingType; import org.apache.pulsar.common.util.ObjectMapperFactory; +import org.apache.pulsar.tests.integration.topologies.PulsarCluster; import org.awaitility.Awaitility; import org.elasticsearch.client.RestClient; import org.elasticsearch.client.RestClientBuilder; @@ -100,6 +99,19 @@ public ElasticSearchSinkTester(boolean schemaEnable) { } } + @Override + protected final ElasticsearchContainer createSinkService(PulsarCluster cluster) { + ElasticsearchContainer elasticContainer = createElasticContainer(); + configureElasticContainer(elasticContainer); + return elasticContainer; + } + + protected void configureElasticContainer(ElasticsearchContainer elasticContainer) { + elasticContainer.withEnv("ingest.geoip.downloader.enabled", "false") + .withLogConsumer(o -> log.info("elastic> {}", o.getUtf8String())); + } + + protected abstract ElasticsearchContainer createElasticContainer(); @Override public void prepareSink() throws Exception { diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/OpenSearchSinkTester.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/OpenSearchSinkTester.java index 1e10cc4189c1a..6f687ca06e7ce 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/OpenSearchSinkTester.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/OpenSearchSinkTester.java @@ -18,9 +18,10 @@ */ package org.apache.pulsar.tests.integration.io.sinks; +import static org.testng.Assert.assertTrue; +import java.util.Map; import java.util.Optional; import org.apache.http.HttpHost; -import org.apache.pulsar.tests.integration.topologies.PulsarCluster; import org.awaitility.Awaitility; import org.opensearch.action.search.SearchRequest; import org.opensearch.action.search.SearchResponse; @@ -31,10 +32,6 @@ import org.testcontainers.elasticsearch.ElasticsearchContainer; import org.testcontainers.utility.DockerImageName; -import java.util.Map; - -import static org.testng.Assert.assertTrue; - public class OpenSearchSinkTester extends ElasticSearchSinkTester { public static final String OPENSEARCH = Optional.ofNullable(System.getenv("OPENSEARCH_IMAGE")) @@ -48,7 +45,7 @@ public OpenSearchSinkTester(boolean schemaEnable) { } @Override - protected ElasticsearchContainer createSinkService(PulsarCluster cluster) { + protected ElasticsearchContainer createElasticContainer() { DockerImageName dockerImageName = DockerImageName.parse(OPENSEARCH) .asCompatibleSubstituteFor("docker.elastic.co/elasticsearch/elasticsearch"); return new ElasticsearchContainer(dockerImageName) From d4363142bfd5e8323ba0229232d8a16447ac0996 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Wed, 28 Jun 2023 11:46:46 +0300 Subject: [PATCH 2/5] [improve][test] Upgrade Elastic 8 container version to 8.5.3 --- .../apache/pulsar/io/elasticsearch/ElasticSearchTestBase.java | 2 +- .../tests/integration/io/sinks/ElasticSearch8SinkTester.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/pulsar-io/elastic-search/src/test/java/org/apache/pulsar/io/elasticsearch/ElasticSearchTestBase.java b/pulsar-io/elastic-search/src/test/java/org/apache/pulsar/io/elasticsearch/ElasticSearchTestBase.java index 9e50a763cf061..a2f5ff380c796 100644 --- a/pulsar-io/elastic-search/src/test/java/org/apache/pulsar/io/elasticsearch/ElasticSearchTestBase.java +++ b/pulsar-io/elastic-search/src/test/java/org/apache/pulsar/io/elasticsearch/ElasticSearchTestBase.java @@ -40,7 +40,7 @@ public abstract class ElasticSearchTestBase { public static final String ELASTICSEARCH_8 = Optional.ofNullable(System.getenv("ELASTICSEARCH_IMAGE_V8")) - .orElse("docker.elastic.co/elasticsearch/elasticsearch:8.5.1"); + .orElse("docker.elastic.co/elasticsearch/elasticsearch:8.5.3"); public static final String ELASTICSEARCH_7 = Optional.ofNullable(System.getenv("ELASTICSEARCH_IMAGE_V7")) .orElse("docker.elastic.co/elasticsearch/elasticsearch:7.17.7"); diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearch8SinkTester.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearch8SinkTester.java index 6ea7d3c1246f8..8e7617a82a5b9 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearch8SinkTester.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearch8SinkTester.java @@ -24,7 +24,7 @@ public class ElasticSearch8SinkTester extends ElasticSearchSinkTester { public static final String ELASTICSEARCH_8 = Optional.ofNullable(System.getenv("ELASTICSEARCH_IMAGE_V8")) - .orElse("docker.elastic.co/elasticsearch/elasticsearch:8.5.1"); + .orElse("docker.elastic.co/elasticsearch/elasticsearch:8.5.3"); public ElasticSearch8SinkTester(boolean schemaEnable) { From ea5cbe812cc19a663e5e87b5c1ab37b851e38ae9 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Wed, 28 Jun 2023 16:34:40 +0300 Subject: [PATCH 3/5] Don't set "ingest.geoip.downloader.enabled" for opensearch container --- .../pulsar/io/elasticsearch/ElasticSearchTestBase.java | 6 ++++-- .../integration/io/sinks/ElasticSearchSinkTester.java | 10 ++++++++-- .../integration/io/sinks/OpenSearchSinkTester.java | 4 ++++ 3 files changed, 16 insertions(+), 4 deletions(-) diff --git a/pulsar-io/elastic-search/src/test/java/org/apache/pulsar/io/elasticsearch/ElasticSearchTestBase.java b/pulsar-io/elastic-search/src/test/java/org/apache/pulsar/io/elasticsearch/ElasticSearchTestBase.java index a2f5ff380c796..0f5a42051c7d1 100644 --- a/pulsar-io/elastic-search/src/test/java/org/apache/pulsar/io/elasticsearch/ElasticSearchTestBase.java +++ b/pulsar-io/elastic-search/src/test/java/org/apache/pulsar/io/elasticsearch/ElasticSearchTestBase.java @@ -73,8 +73,10 @@ protected ElasticsearchContainer createElasticsearchContainer() { } protected void configureElasticContainer(ElasticsearchContainer elasticContainer) { - elasticContainer.withEnv("ingest.geoip.downloader.enabled", "false") - .withLogConsumer(o -> log.info("elastic> {}", o.getUtf8String())); + if (getCompatibilityMode() != ElasticSearchConfig.CompatibilityMode.OPENSEARCH) { + elasticContainer.withEnv("ingest.geoip.downloader.enabled", "false"); + } + elasticContainer.withLogConsumer(o -> log.info("elastic> {}", o.getUtf8String())); } protected ElasticSearchConfig.CompatibilityMode getCompatibilityMode() { diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearchSinkTester.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearchSinkTester.java index c7d8c84f684a0..65b97ee2d6852 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearchSinkTester.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearchSinkTester.java @@ -107,8 +107,14 @@ protected final ElasticsearchContainer createSinkService(PulsarCluster cluster) } protected void configureElasticContainer(ElasticsearchContainer elasticContainer) { - elasticContainer.withEnv("ingest.geoip.downloader.enabled", "false") - .withLogConsumer(o -> log.info("elastic> {}", o.getUtf8String())); + if (!isOpenSearch()) { + elasticContainer.withEnv("ingest.geoip.downloader.enabled", "false"); + } + elasticContainer.withLogConsumer(o -> log.info("elastic> {}", o.getUtf8String())); + } + + protected boolean isOpenSearch() { + return false; } protected abstract ElasticsearchContainer createElasticContainer(); diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/OpenSearchSinkTester.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/OpenSearchSinkTester.java index 6f687ca06e7ce..75f0fdac6f90c 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/OpenSearchSinkTester.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/OpenSearchSinkTester.java @@ -54,6 +54,10 @@ protected ElasticsearchContainer createElasticContainer() { .withEnv("plugins.security.disabled", "true"); } + protected boolean isOpenSearch() { + return true; + } + @Override public void prepareSink() throws Exception { RestClientBuilder builder = RestClient.builder( From 236307e217931292f18ff13d3785e6fa4869b3fc Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Wed, 28 Jun 2023 18:30:50 +0300 Subject: [PATCH 4/5] Set diskUsageWarnThreshold and diskUsageLwmThreshold for int tests --- .../pulsar/tests/integration/topologies/PulsarCluster.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/topologies/PulsarCluster.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/topologies/PulsarCluster.java index c4c7697e30fdd..9b4823f46d4cc 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/topologies/PulsarCluster.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/topologies/PulsarCluster.java @@ -167,7 +167,9 @@ private PulsarCluster(PulsarClusterSpec spec, CSContainer csContainer, boolean s .withEnv("journalSyncData", "false") .withEnv("journalMaxGroupWaitMSec", "0") .withEnv("clusterName", clusterName) + .withEnv("PULSAR_PREFIX_diskUsageWarnThreshold", "0.95") .withEnv("diskUsageThreshold", "0.99") + .withEnv("PULSAR_PREFIX_diskUsageLwmThreshold", "0.97") .withEnv("nettyMaxFrameSizeBytes", String.valueOf(spec.maxMessageSize)); if (spec.bookkeeperEnvs != null) { bookieContainer.withEnv(spec.bookkeeperEnvs); From d959eb4929d4192fb56c140a8b590e0ba25d866b Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Wed, 28 Jun 2023 18:38:43 +0300 Subject: [PATCH 5/5] Configure Elastic to allow the disk fill beyond 90% - similar setting used in https://github.com/elastic/elastic-github-actions/blob/562b8b6ae4677da97273ff6bc4d630ce96ecbaa5/elasticsearch/run-elasticsearch.sh#L41 --- .../tests/integration/io/sinks/ElasticSearchSinkTester.java | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearchSinkTester.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearchSinkTester.java index 65b97ee2d6852..0784055d29003 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearchSinkTester.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/io/sinks/ElasticSearchSinkTester.java @@ -110,6 +110,10 @@ protected void configureElasticContainer(ElasticsearchContainer elasticContainer if (!isOpenSearch()) { elasticContainer.withEnv("ingest.geoip.downloader.enabled", "false"); } + + // allow disk to fill up beyond default 90% threshold + elasticContainer.withEnv("cluster.routing.allocation.disk.threshold_enabled", "false"); + elasticContainer.withLogConsumer(o -> log.info("elastic> {}", o.getUtf8String())); }