diff --git a/applications/sink/elasticsearch-sink/pom.xml b/applications/sink/elasticsearch-sink/pom.xml
index 83ffd2ab..fa8c3196 100644
--- a/applications/sink/elasticsearch-sink/pom.xml
+++ b/applications/sink/elasticsearch-sink/pom.xml
@@ -33,19 +33,16 @@
org.testcontainers
testcontainers
- ${testcontainers.version}
test
org.testcontainers
junit-jupiter
- ${testcontainers.version}
test
org.testcontainers
elasticsearch
- ${testcontainers.version}
test
diff --git a/applications/sink/elasticsearch-sink/src/test/java/org/springframework/cloud/stream/app/sink/elasticsearch/ElasticsearchSinkTests.java b/applications/sink/elasticsearch-sink/src/test/java/org/springframework/cloud/stream/app/sink/elasticsearch/ElasticsearchSinkTests.java
index d3bf671b..e60a246f 100644
--- a/applications/sink/elasticsearch-sink/src/test/java/org/springframework/cloud/stream/app/sink/elasticsearch/ElasticsearchSinkTests.java
+++ b/applications/sink/elasticsearch-sink/src/test/java/org/springframework/cloud/stream/app/sink/elasticsearch/ElasticsearchSinkTests.java
@@ -32,14 +32,8 @@ import org.testcontainers.utility.DockerImageName;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
-import org.springframework.cloud.fn.consumer.elasticsearch.ElasticsearchConsumerConfiguration;
import org.springframework.cloud.stream.binder.test.InputDestination;
import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
-import org.springframework.context.annotation.Configuration;
-import org.springframework.context.annotation.Import;
-import org.springframework.data.elasticsearch.client.ClientConfiguration;
-import org.springframework.data.elasticsearch.client.elc.ElasticsearchConfiguration;
-import org.springframework.lang.NonNull;
import org.springframework.messaging.support.GenericMessage;
@@ -59,24 +53,20 @@ public class ElasticsearchSinkTests {
private final ApplicationContextRunner contextRunner = new ApplicationContextRunner()
.withUserConfiguration(TestChannelBinderConfiguration.getCompleteConfiguration(ElasticsearchSinkTestApplication.class));
-
@Test
void elasticSearchSinkWithIndexNameProperty() {
this.contextRunner
.withPropertyValues("spring.cloud.function.definition=elasticsearchConsumer",
"elasticsearch.consumer.index=foo", "elasticsearch.consumer.id=1",
- "spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
+ "spring.elasticsearch.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
-
- final InputDestination inputDestination = context.getBean(InputDestination.class);
- final String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564,"
+ InputDestination inputDestination = context.getBean(InputDestination.class);
+ String jsonObject = "{\"age\":10,\"dateOfBirth\":1471466076564,"
+ "\"fullName\":\"John Doe\"}";
-
inputDestination.send(new GenericMessage<>(jsonObject));
-
- final ElasticsearchClient elasticsearchClient = context.getBean(ElasticsearchClient.class);
- final GetRequest getRequest = new GetRequest.Builder().index("foo").id("1").build();
- final GetResponse response = elasticsearchClient.get(getRequest, JsonData.class);
+ ElasticsearchClient elasticsearchClient = context.getBean(ElasticsearchClient.class);
+ GetRequest getRequest = new GetRequest.Builder().index("foo").id("1").build();
+ GetResponse response = elasticsearchClient.get(getRequest, JsonData.class);
assertThat(response.found()).isTrue();
assertThat(response.source()).isNotNull();
assertThat(response.source().toJson()).isEqualTo(JsonData.fromJson(jsonObject).toJson());
@@ -87,7 +77,7 @@ public class ElasticsearchSinkTests {
void elasticSearchSinkWithIndexNameFromHeader() {
this.contextRunner
.withPropertyValues("spring.cloud.function.definition=elasticsearchConsumer", "elasticsearch.consumer.id=1",
- "spring.elasticsearch.rest.uris=http://" + elasticsearch.getHttpHostAddress())
+ "spring.elasticsearch.uris=http://" + elasticsearch.getHttpHostAddress())
.run(context -> {
final InputDestination inputDestination = context.getBean(InputDestination.class);
@@ -106,19 +96,7 @@ public class ElasticsearchSinkTests {
}
@SpringBootApplication
- @Import(ElasticsearchConsumerConfiguration.class)
static class ElasticsearchSinkTestApplication {
}
- @Configuration(proxyBeanMethods = false)
- static class Config extends ElasticsearchConfiguration {
- @NonNull
- @Override
- public ClientConfiguration clientConfiguration() {
- return ClientConfiguration.builder()
- .connectedTo(elasticsearch.getHttpHostAddress())
- .build();
- }
- }
-
}
diff --git a/applications/sink/mqtt-sink/pom.xml b/applications/sink/mqtt-sink/pom.xml
index 80c44c21..8c987572 100644
--- a/applications/sink/mqtt-sink/pom.xml
+++ b/applications/sink/mqtt-sink/pom.xml
@@ -22,7 +22,6 @@
org.testcontainers
testcontainers
- ${testcontainers.version}
test
diff --git a/applications/sink/redis-sink/pom.xml b/applications/sink/redis-sink/pom.xml
index b99a65c1..47235808 100644
--- a/applications/sink/redis-sink/pom.xml
+++ b/applications/sink/redis-sink/pom.xml
@@ -23,19 +23,16 @@
org.testcontainers
testcontainers
- ${testcontainers.version}
test
org.testcontainers
junit-jupiter
- ${testcontainers.version}
test
org.springframework.cloud
spring-cloud-starter-bootstrap
- ${spring-cloud-starters.version}
test
diff --git a/applications/source/debezium-source/src/test/java/org/springframework/cloud/stream/app/source/debezium/integration/DebeziumSupplierAvroFormatTest.java b/applications/source/debezium-source/src/test/java/org/springframework/cloud/stream/app/source/debezium/integration/DebeziumSupplierAvroFormatTest.java
deleted file mode 100644
index fe8df20e..00000000
--- a/applications/source/debezium-source/src/test/java/org/springframework/cloud/stream/app/source/debezium/integration/DebeziumSupplierAvroFormatTest.java
+++ /dev/null
@@ -1,121 +0,0 @@
-/*
- * Copyright 2023-2023 the original author or authors.
- *
- * Licensed under the Apache License, Version 2.0 (the "License");
- * you may not use this file except in compliance with the License.
- * You may obtain a copy of the License at
- *
- * https://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.springframework.cloud.stream.app.source.debezium.integration;
-
-import java.time.Duration;
-import java.util.List;
-
-import org.junit.jupiter.api.Tag;
-import org.junit.jupiter.api.Test;
-import org.testcontainers.containers.GenericContainer;
-import org.testcontainers.junit.jupiter.Container;
-import org.testcontainers.junit.jupiter.Testcontainers;
-
-import org.springframework.boot.WebApplicationType;
-import org.springframework.boot.builder.SpringApplicationBuilder;
-import org.springframework.cloud.stream.binder.test.OutputDestination;
-import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
-import org.springframework.context.ConfigurableApplicationContext;
-import org.springframework.messaging.Message;
-
-import static org.assertj.core.api.Assertions.assertThat;
-
-/**
- * @author Christian Tzolov
- */
-@Tag("integration")
-@Testcontainers
-public class DebeziumSupplierAvroFormatTest {
-
- // E.g. docker run -it --rm --name apicurio -p 8080:8080 apicurio/apicurio-registry-mem:2.4.1.Final
- @Container
- static GenericContainer> apicurio = new GenericContainer<>("apicurio/apicurio-registry-mem:2.4.1.Final")
- .withExposedPorts(8080)
- .withStartupTimeout(Duration.ofSeconds(120))
- .withStartupAttempts(3);
-
- @Container
- static GenericContainer> debeziumMySQL = new GenericContainer<>(DebeziumTestUtils.DEBEZIUM_EXAMPLE_MYSQL_IMAGE)
- .withEnv("MYSQL_ROOT_PASSWORD", "debezium")
- .withEnv("MYSQL_USER", "mysqluser")
- .withEnv("MYSQL_PASSWORD", "mysqlpw")
- .withExposedPorts(3306)
- .withStartupTimeout(Duration.ofSeconds(120))
- .withStartupAttempts(3);
-
- private final SpringApplicationBuilder applicationBuilder = new SpringApplicationBuilder(
- TestChannelBinderConfiguration.getCompleteConfiguration(TestDebeziumSourceApplication.class))
- .web(WebApplicationType.NONE)
- .properties(
- "spring.cloud.function.definition=debeziumSupplier",
-
- "debezium.payloadFormat=AVRO",
-
- "debezium.properties.key.converter=io.apicurio.registry.utils.converter.AvroConverter",
- "debezium.properties.key.converter.apicurio.registry.auto-register=true",
- "debezium.properties.key.converter.apicurio.registry.find-latest=true",
- "debezium.properties.value.converter=io.apicurio.registry.utils.converter.AvroConverter",
- "debezium.properties.value.converter.apicurio.registry.auto-register=true",
- "debezium.properties.value.converter.apicurio.registry.find-latest=true",
- "debezium.properties.schema.name.adjustment.mode=avro",
-
- "debezium.properties.schema.history.internal=io.debezium.relational.history.MemorySchemaHistory",
- "debezium.properties.offset.storage=org.apache.kafka.connect.storage.MemoryOffsetBackingStore",
-
- "debezium.properties.topic.prefix=my-topic",
- "debezium.properties.name=my-connector",
- "debezium.properties.database.server.id=85744",
-
- "debezium.properties.connector.class=io.debezium.connector.mysql.MySqlConnector",
- "debezium.properties.database.user=debezium",
- "debezium.properties.database.password=dbz",
- "debezium.properties.database.hostname=localhost",
-
- // JdbcTemplate configuration
- String.format("app.datasource.url=jdbc:mysql://localhost:%d/%s?enabledTLSProtocols=TLSv1.2",
- debeziumMySQL.getMappedPort(3306), DebeziumTestUtils.DATABASE_NAME),
- "app.datasource.username=root",
- "app.datasource.password=debezium",
- "app.datasource.driver-class-name=com.mysql.cj.jdbc.Driver",
- "app.datasource.type=com.zaxxer.hikari.HikariDataSource");
-
- @Test
- public void mysqlWithAvroContentFormat() {
-
- String MYSQL_MAPPED_PORT = String.valueOf(debeziumMySQL.getMappedPort(3306));
- String APICURIO_URL = "http://localhost:" + String.valueOf(apicurio.getMappedPort(8080)) + "/apis/registry/v2";
-
- try (ConfigurableApplicationContext context = applicationBuilder.run(
- "--debezium.properties.key.converter.apicurio.registry.url=" + APICURIO_URL,
- "--debezium.properties.value.converter.apicurio.registry.url=" + APICURIO_URL,
- "--debezium.properties.database.port=" + MYSQL_MAPPED_PORT)) {
-
- OutputDestination outputDestination = context.getBean(OutputDestination.class);
-
- // Using local region here
- List> messages = DebeziumTestUtils.receiveAll(outputDestination);
-
- assertThat(messages).isNotNull();
- // Message size should correspond to the number of insert statements in the sample inventor DB
- // configured by:
- // https://github.com/debezium/container-images/blob/main/examples/mysql/2.1/inventory.sql
- assertThat(messages).hasSizeGreaterThanOrEqualTo(52);
- // assertThat(messages).map(message ->
- // message.getHeaders().get("contentType")).isEqualTo("application/avro"); // TEST utils bug.
- }
- }
-}
diff --git a/applications/source/mqtt-source/pom.xml b/applications/source/mqtt-source/pom.xml
index d63afa6f..6fdc46fc 100644
--- a/applications/source/mqtt-source/pom.xml
+++ b/applications/source/mqtt-source/pom.xml
@@ -22,7 +22,6 @@
org.testcontainers
testcontainers
- ${testcontainers.version}
test
diff --git a/applications/source/websocket-source/src/test/java/org/springframework/cloud/stream/app/source/websocket/WebsocketSourceTests.java b/applications/source/websocket-source/src/test/java/org/springframework/cloud/stream/app/source/websocket/WebsocketSourceTests.java
index c9472142..e7f9fd0d 100644
--- a/applications/source/websocket-source/src/test/java/org/springframework/cloud/stream/app/source/websocket/WebsocketSourceTests.java
+++ b/applications/source/websocket-source/src/test/java/org/springframework/cloud/stream/app/source/websocket/WebsocketSourceTests.java
@@ -18,6 +18,7 @@ package org.springframework.cloud.stream.app.source.websocket;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
+import java.util.Base64;
import java.util.function.Supplier;
import org.junit.jupiter.api.Test;
@@ -37,7 +38,6 @@ import org.springframework.http.HttpHeaders;
import org.springframework.integration.websocket.ClientWebSocketContainer;
import org.springframework.messaging.Message;
import org.springframework.test.annotation.DirtiesContext;
-import org.springframework.util.Base64Utils;
import org.springframework.web.socket.TextMessage;
import org.springframework.web.socket.WebSocketSession;
import org.springframework.web.socket.client.standard.StandardWebSocketClient;
@@ -73,7 +73,7 @@ public class WebsocketSourceTests {
this.properties.getPath());
HttpHeaders httpHeaders = new HttpHeaders();
- String token = Base64Utils.encodeToString(
+ String token = Base64.getEncoder().encodeToString(
(this.securityProperties.getUser().getName() + ":" + this.securityProperties.getUser().getPassword())
.getBytes(StandardCharsets.UTF_8));
httpHeaders.set(HttpHeaders.AUTHORIZATION, "Basic " + token);
diff --git a/applications/stream-applications-core/pom.xml b/applications/stream-applications-core/pom.xml
index c5a578d5..59e569e1 100644
--- a/applications/stream-applications-core/pom.xml
+++ b/applications/stream-applications-core/pom.xml
@@ -17,7 +17,7 @@
6.0.0-SNAPSHOT
- 5.0.1
+ 5.1.0-SNAPSHOT
1.5.3
1.1.0-SNAPSHOT
1.1.0-SNAPSHOT
@@ -42,38 +42,12 @@
${java-functions.version}
pom
import
-
-
- org.springframework.cloud
- spring-cloud-function-dependencies
- ${spring-cloud-function.version}
-
-
- org.springframework.cloud
- spring-cloud-stream-dependencies
- ${spring-cloud-stream-dependencies.version}
- pom
- import
org.mock-server
mockserver-netty
${mockserver.version}
-
- org.junit
- junit-bom
- 5.9.3
- pom
- import
-
-
- org.testcontainers
- testcontainers-bom
- ${testcontainers.version}
- pom
- import
-
@@ -169,17 +143,6 @@
-
-
- org.springframework.cloud
- spring-cloud-function-dependencies
- ${spring-cloud-function.version}
-
-
- org.springframework.cloud
- spring-cloud-stream-dependencies
- ${spring-cloud-stream-dependencies.version}
-
org.springframework.cloud
spring-cloud-dependencies
diff --git a/applications/stream-applications-integration-tests/pom.xml b/applications/stream-applications-integration-tests/pom.xml
index bcd497d5..10d02cbd 100644
--- a/applications/stream-applications-integration-tests/pom.xml
+++ b/applications/stream-applications-integration-tests/pom.xml
@@ -101,7 +101,6 @@
org.testcontainers
localstack
- ${testcontainers.version}
diff --git a/stream-applications-build/pom.xml b/stream-applications-build/pom.xml
index 173524c5..101356e1 100644
--- a/stream-applications-build/pom.xml
+++ b/stream-applications-build/pom.xml
@@ -35,26 +35,26 @@
0.0.11
true
0.0.43
- 3.3.6
- 6.1.15
- 2023.0.4
- 4.1.5
- 4.1.4
- 4.1.4
- 1.19.8
+ 3.4.2-SNAPSHOT
+ 2024.0.1-SNAPSHOT
+
+
+
+ 6.2.1
+ 4.2.1-SNAPSHOT
+
5.15.0
- 4.0.24
2.5.2
- 3.1.0
5.5.0
1
1.26.2
3.17.0
-
- ${env.BUILD_NAME}
-
- ${env.BUILD_NUMBER}
-
+ 4.12.0
+
+ ${env.BUILD_NAME}
+
+ ${env.BUILD_NUMBER}
+
https://spring.io/projects/spring-cloud-stream-applications
@@ -82,12 +82,6 @@
commons-lang3
${commons-lang3.version}
-
- org.springframework
- spring-framework-bom
- ${spring.version}
- pom
-
org.springframework.boot
spring-boot-dependencies
@@ -102,37 +96,18 @@
pom
import
-
- org.springframework.cloud
- spring-cloud-function-dependencies
- ${spring-cloud-function.version}
- import
- pom
-
-
- org.testcontainers
- testcontainers-bom
- ${testcontainers.version}
- pom
- import
-
-
- org.apache.groovy
- groovy-bom
- ${groovy.version}
- pom
- import
-
-
- jakarta.jms
- jakarta.jms-api
- ${jakarta-jms.version}
-
org.apache.ivy
ivy
${apache-ivy.version}
+
+ com.squareup.okhttp3
+ okhttp-bom
+ ${okhttp3.version}
+ pom
+ import
+