@@ -21,7 +21,7 @@ ext {
|
||||
micrometerTracingVersion = '1.0.0-SNAPSHOT'
|
||||
mockitoVersion = '4.6.1'
|
||||
protobufJavaVersion = '3.21.5'
|
||||
pulsarTestcontainersVersion = '1.17.3'
|
||||
testcontainersVersion = '1.17.3'
|
||||
pulsarVersion = '2.10.1'
|
||||
reactorVersion = '2022.0.0-SNAPSHOT'
|
||||
springBootVersion = '3.0.0-SNAPSHOT'
|
||||
@@ -37,6 +37,7 @@ dependencies {
|
||||
api platform("org.junit:junit-bom:$junitJupiterVersion")
|
||||
api platform("org.mockito:mockito-bom:$mockitoVersion")
|
||||
api platform("org.springframework:spring-framework-bom:$springVersion")
|
||||
api platform("org.testcontainers:testcontainers-bom:$testcontainersVersion")
|
||||
api platform("io.micrometer:micrometer-bom:$micrometerVersion")
|
||||
api platform("io.micrometer:micrometer-tracing-bom:$micrometerTracingVersion")
|
||||
|
||||
@@ -58,6 +59,5 @@ dependencies {
|
||||
api "org.springframework.boot:spring-boot-starter-validation:$springBootVersion"
|
||||
api "org.springframework.boot:spring-boot-starter-test:$springBootVersion"
|
||||
api "org.springframework.retry:spring-retry:$springRetryVersion"
|
||||
api "org.testcontainers:pulsar:$pulsarTestcontainersVersion"
|
||||
}
|
||||
}
|
||||
|
||||
@@ -19,5 +19,6 @@ dependencies {
|
||||
testRuntimeOnly 'org.apache.logging.log4j:log4j-core'
|
||||
testRuntimeOnly 'org.apache.logging.log4j:log4j-jcl'
|
||||
testImplementation 'org.springframework.boot:spring-boot-starter-test'
|
||||
testImplementation 'org.testcontainers:junit-jupiter'
|
||||
testImplementation 'org.testcontainers:pulsar'
|
||||
}
|
||||
|
||||
@@ -38,7 +38,7 @@ import org.springframework.pulsar.core.PulsarTemplate;
|
||||
* @author Soby Chacko
|
||||
* @author Chris Bono
|
||||
*/
|
||||
class PulsarListenerTests extends AbstractContainerBaseTests {
|
||||
class PulsarListenerTests implements PulsarTestContainerSupport {
|
||||
|
||||
static CountDownLatch latch1 = new CountDownLatch(1);
|
||||
static CountDownLatch latch2 = new CountDownLatch(10);
|
||||
@@ -49,7 +49,7 @@ class PulsarListenerTests extends AbstractContainerBaseTests {
|
||||
app.setWebApplicationType(WebApplicationType.NONE);
|
||||
|
||||
try (ConfigurableApplicationContext context = app
|
||||
.run("--spring.pulsar.client.serviceUrl=" + AbstractContainerBaseTests.getPulsarBrokerUrl())) {
|
||||
.run("--spring.pulsar.client.serviceUrl=" + PulsarTestContainerSupport.getPulsarBrokerUrl())) {
|
||||
@SuppressWarnings("unchecked")
|
||||
final PulsarTemplate<String> pulsarTemplate = context.getBean(PulsarTemplate.class);
|
||||
pulsarTemplate.send("hello-pulsar-exclusive", "John Doe");
|
||||
@@ -64,7 +64,7 @@ class PulsarListenerTests extends AbstractContainerBaseTests {
|
||||
app.setWebApplicationType(WebApplicationType.NONE);
|
||||
|
||||
try (ConfigurableApplicationContext context = app
|
||||
.run("--spring.pulsar.client.serviceUrl=" + AbstractContainerBaseTests.getPulsarBrokerUrl())) {
|
||||
.run("--spring.pulsar.client.serviceUrl=" + PulsarTestContainerSupport.getPulsarBrokerUrl())) {
|
||||
@SuppressWarnings("unchecked")
|
||||
final PulsarTemplate<String> pulsarTemplate = context.getBean(PulsarTemplate.class);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
@@ -77,7 +77,7 @@ class PulsarListenerTests extends AbstractContainerBaseTests {
|
||||
|
||||
@Configuration
|
||||
@Import(PulsarAutoConfiguration.class)
|
||||
public static class BasicListenerConfig {
|
||||
static class BasicListenerConfig {
|
||||
|
||||
@PulsarListener(subscriptionName = "test-exclusive-sub-1", topics = "hello-pulsar-exclusive")
|
||||
public void listen(String foo) {
|
||||
@@ -88,7 +88,7 @@ class PulsarListenerTests extends AbstractContainerBaseTests {
|
||||
|
||||
@Configuration
|
||||
@Import(PulsarAutoConfiguration.class)
|
||||
public static class BatchListenerConfig {
|
||||
static class BatchListenerConfig {
|
||||
|
||||
@PulsarListener(subscriptionName = "test-exclusive-sub-2", topics = "hello-pulsar-exclusive", batch = true)
|
||||
public void listen(List<String> foo) {
|
||||
|
||||
@@ -18,24 +18,32 @@ package org.springframework.pulsar.autoconfigure;
|
||||
|
||||
import java.util.Locale;
|
||||
|
||||
import org.junit.jupiter.api.BeforeAll;
|
||||
import org.testcontainers.containers.PulsarContainer;
|
||||
import org.testcontainers.junit.jupiter.Testcontainers;
|
||||
import org.testcontainers.utility.DockerImageName;
|
||||
|
||||
abstract class AbstractContainerBaseTests {
|
||||
/**
|
||||
* Provides a static {@link PulsarContainer} that can be shared across test classes.
|
||||
*
|
||||
* @author Chris Bono
|
||||
*/
|
||||
@Testcontainers(disabledWithoutDocker = true)
|
||||
public interface PulsarTestContainerSupport {
|
||||
|
||||
static final PulsarContainer PULSAR_CONTAINER;
|
||||
PulsarContainer PULSAR_CONTAINER = new PulsarContainer(
|
||||
isRunningOnMacM1() ? getMacM1PulsarImage() : getStandardPulsarImage());
|
||||
|
||||
static {
|
||||
final DockerImageName PULSAR_IMAGE = isRunningOnMacM1() ? getMacM1PulsarImage() : getStandardPulsarImage();
|
||||
PULSAR_CONTAINER = new PulsarContainer(PULSAR_IMAGE);
|
||||
@BeforeAll
|
||||
static void startContainer() {
|
||||
PULSAR_CONTAINER.start();
|
||||
}
|
||||
|
||||
protected static String getPulsarBrokerUrl() {
|
||||
static String getPulsarBrokerUrl() {
|
||||
return PULSAR_CONTAINER.getPulsarBrokerUrl();
|
||||
}
|
||||
|
||||
protected static String getHttpServiceUrl() {
|
||||
static String getHttpServiceUrl() {
|
||||
return PULSAR_CONTAINER.getHttpServiceUrl();
|
||||
}
|
||||
|
||||
@@ -35,7 +35,7 @@ import org.springframework.pulsar.core.PulsarTemplate;
|
||||
* @author Chris Bono
|
||||
*/
|
||||
@SpringBootTest(classes = SpringPulsarBootTestApp.class)
|
||||
public class SpringPulsarBootAppSanityTests extends AbstractContainerBaseTests {
|
||||
class SpringPulsarBootAppSanityTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
void appStartsWithAutoConfiguredSpringPulsarComponents(
|
||||
|
||||
@@ -4,7 +4,7 @@
|
||||
<pattern>%d{HH:mm:ss.SSS} [%thread] %-5level %logger - %msg%n</pattern>
|
||||
</encoder>
|
||||
</appender>
|
||||
<root level="info">
|
||||
<root level="WARN">
|
||||
<appender-ref ref="STDOUT"/>
|
||||
</root>
|
||||
<logger name="org.testcontainers" level="ERROR"/>
|
||||
|
||||
@@ -36,6 +36,7 @@ dependencies {
|
||||
testImplementation 'org.hamcrest:hamcrest'
|
||||
testImplementation 'org.mockito:mockito-junit-jupiter'
|
||||
testImplementation 'org.springframework:spring-test'
|
||||
testImplementation 'org.testcontainers:junit-jupiter'
|
||||
testImplementation 'org.testcontainers:pulsar'
|
||||
}
|
||||
|
||||
|
||||
@@ -58,7 +58,7 @@ import org.springframework.pulsar.listener.PulsarRecordMessageListener;
|
||||
/**
|
||||
* @author Soby Chacko
|
||||
*/
|
||||
class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
|
||||
class ConsumerAcknowledgmentTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
void testRecordAck() throws Exception {
|
||||
@@ -67,7 +67,8 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
|
||||
strings.add("cons-ack-tests-011");
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "cons-ack-tests-sb-011");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
|
||||
@@ -108,7 +109,8 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
|
||||
strings.add("cons-ack-tests-012");
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "cons-ack-tests-sb-012");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
|
||||
@@ -144,7 +146,8 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
|
||||
Map<String, Object> config = new HashMap<>();
|
||||
config.put("topicNames", Collections.singleton("cons-ack-tests-013"));
|
||||
config.put("subscriptionName", "cons-ack-tests-sb-013");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
|
||||
@@ -212,7 +215,8 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
|
||||
strings.add("cons-ack-tests-014");
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "cons-ack-tests-sb-014");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
|
||||
@@ -266,7 +270,8 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
|
||||
strings.add("cons-ack-tests-015");
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "cons-ack-tests-sb-015");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
|
||||
@@ -314,7 +319,8 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
|
||||
strings.add("cons-ack-tests-016");
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "cons-ack-tests-sb-016");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
|
||||
|
||||
@@ -42,11 +42,12 @@ import org.springframework.pulsar.listener.PulsarRecordMessageListener;
|
||||
/**
|
||||
* @author Soby Chacko
|
||||
*/
|
||||
class FailoverConsumerTests extends AbstractContainerBaseTests {
|
||||
class FailoverConsumerTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
void testFailOverConsumersOnPartitionedTopic() throws Exception {
|
||||
PulsarAdmin admin = PulsarAdmin.builder().serviceHttpUrl(getHttpServiceUrl()).build();
|
||||
PulsarAdmin admin = PulsarAdmin.builder().serviceHttpUrl(PulsarTestContainerSupport.getHttpServiceUrl())
|
||||
.build();
|
||||
|
||||
String topicName = "persistent://public/default/my-part-topic-1";
|
||||
int numPartitions = 3;
|
||||
@@ -57,7 +58,8 @@ class FailoverConsumerTests extends AbstractContainerBaseTests {
|
||||
topics.add("my-part-topic-1");
|
||||
config.put("topicNames", topics);
|
||||
config.put("subscriptionName", "my-part-subscription-1");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
CountDownLatch latch = new CountDownLatch(3);
|
||||
|
||||
@@ -43,7 +43,7 @@ import org.springframework.test.context.junit.jupiter.SpringExtension;
|
||||
*/
|
||||
@ExtendWith(SpringExtension.class)
|
||||
@ContextConfiguration
|
||||
public class PulsarAdministrationTests extends AbstractContainerBaseTests {
|
||||
public class PulsarAdministrationTests implements PulsarTestContainerSupport {
|
||||
|
||||
private static final String NAMESPACE = "public/default";
|
||||
|
||||
@@ -73,12 +73,13 @@ public class PulsarAdministrationTests extends AbstractContainerBaseTests {
|
||||
|
||||
@Bean
|
||||
PulsarAdmin pulsarAdminClient() throws PulsarClientException {
|
||||
return PulsarAdmin.builder().serviceHttpUrl(getHttpServiceUrl()).build();
|
||||
return PulsarAdmin.builder().serviceHttpUrl(PulsarTestContainerSupport.getHttpServiceUrl()).build();
|
||||
}
|
||||
|
||||
@Bean
|
||||
PulsarAdministration pulsarAdministration() {
|
||||
return new PulsarAdministration(PulsarAdmin.builder().serviceHttpUrl(getHttpServiceUrl()));
|
||||
return new PulsarAdministration(
|
||||
PulsarAdmin.builder().serviceHttpUrl(PulsarTestContainerSupport.getHttpServiceUrl()));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -48,7 +48,7 @@ import org.junit.jupiter.api.Test;
|
||||
* @author Chris Bono
|
||||
* @author Alexander Preuß
|
||||
*/
|
||||
abstract class PulsarProducerFactoryTests extends AbstractContainerBaseTests {
|
||||
abstract class PulsarProducerFactoryTests implements PulsarTestContainerSupport {
|
||||
|
||||
protected final Schema<String> schema = Schema.STRING;
|
||||
|
||||
@@ -56,7 +56,7 @@ abstract class PulsarProducerFactoryTests extends AbstractContainerBaseTests {
|
||||
|
||||
@BeforeEach
|
||||
void createPulsarClient() throws PulsarClientException {
|
||||
pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
}
|
||||
|
||||
@AfterEach
|
||||
|
||||
@@ -59,12 +59,13 @@ import org.springframework.pulsar.core.PulsarOperations.SendMessageBuilder;
|
||||
* @author Chris Bono
|
||||
* @author Alexander Preuß
|
||||
*/
|
||||
class PulsarTemplateTests extends AbstractContainerBaseTests {
|
||||
class PulsarTemplateTests implements PulsarTestContainerSupport {
|
||||
|
||||
@ParameterizedTest(name = "{0}")
|
||||
@MethodSource("interceptorInvocationTestProvider")
|
||||
void interceptorInvocationTest(String topic, List<ProducerInterceptor> interceptors) throws Exception {
|
||||
try (PulsarClient client = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build()) {
|
||||
try (PulsarClient client = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build()) {
|
||||
PulsarProducerFactory<String> producerFactory = new DefaultPulsarProducerFactory<>(client,
|
||||
Collections.singletonMap("topicName", topic));
|
||||
PulsarTemplate<String> pulsarTemplate = new PulsarTemplate<>(producerFactory, interceptors);
|
||||
@@ -86,7 +87,8 @@ class PulsarTemplateTests extends AbstractContainerBaseTests {
|
||||
@Test
|
||||
void sendMessageWithSpecificSchemaTest() throws Exception {
|
||||
String topic = "smt-specific-schema-topic";
|
||||
try (PulsarClient client = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build()) {
|
||||
try (PulsarClient client = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build()) {
|
||||
try (Consumer<Foo> consumer = client.newConsumer(Schema.JSON(Foo.class)).topic(topic)
|
||||
.subscriptionName("test-specific-schema-subscription").subscribe()) {
|
||||
PulsarProducerFactory<Foo> producerFactory = new DefaultPulsarProducerFactory<>(client,
|
||||
@@ -123,11 +125,13 @@ class PulsarTemplateTests extends AbstractContainerBaseTests {
|
||||
}
|
||||
|
||||
if (router != null) {
|
||||
try (PulsarAdmin admin = PulsarAdmin.builder().serviceHttpUrl(getHttpServiceUrl()).build()) {
|
||||
try (PulsarAdmin admin = PulsarAdmin.builder()
|
||||
.serviceHttpUrl(PulsarTestContainerSupport.getHttpServiceUrl()).build()) {
|
||||
admin.topics().createPartitionedTopic("persistent://public/default/" + topic, 1);
|
||||
}
|
||||
}
|
||||
try (PulsarClient client = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build()) {
|
||||
try (PulsarClient client = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl())
|
||||
.build()) {
|
||||
try (Consumer<String> consumer = client.newConsumer(Schema.STRING).topic(topic)
|
||||
.subscriptionName(subscription).subscribe()) {
|
||||
Map<String, Object> producerConfig = testArgs.useSpecificTopic ? Collections.emptyMap()
|
||||
|
||||
@@ -18,24 +18,32 @@ package org.springframework.pulsar.core;
|
||||
|
||||
import java.util.Locale;
|
||||
|
||||
import org.junit.jupiter.api.BeforeAll;
|
||||
import org.testcontainers.containers.PulsarContainer;
|
||||
import org.testcontainers.junit.jupiter.Testcontainers;
|
||||
import org.testcontainers.utility.DockerImageName;
|
||||
|
||||
public abstract class AbstractContainerBaseTests {
|
||||
/**
|
||||
* Provides a static {@link PulsarContainer} that can be shared across test classes.
|
||||
*
|
||||
* @author Chris Bono
|
||||
*/
|
||||
@Testcontainers(disabledWithoutDocker = true)
|
||||
public interface PulsarTestContainerSupport {
|
||||
|
||||
static final PulsarContainer PULSAR_CONTAINER;
|
||||
PulsarContainer PULSAR_CONTAINER = new PulsarContainer(
|
||||
isRunningOnMacM1() ? getMacM1PulsarImage() : getStandardPulsarImage());
|
||||
|
||||
static {
|
||||
final DockerImageName PULSAR_IMAGE = isRunningOnMacM1() ? getMacM1PulsarImage() : getStandardPulsarImage();
|
||||
PULSAR_CONTAINER = new PulsarContainer(PULSAR_IMAGE);
|
||||
@BeforeAll
|
||||
static void startContainer() {
|
||||
PULSAR_CONTAINER.start();
|
||||
}
|
||||
|
||||
protected static String getPulsarBrokerUrl() {
|
||||
static String getPulsarBrokerUrl() {
|
||||
return PULSAR_CONTAINER.getPulsarBrokerUrl();
|
||||
}
|
||||
|
||||
protected static String getHttpServiceUrl() {
|
||||
static String getHttpServiceUrl() {
|
||||
return PULSAR_CONTAINER.getHttpServiceUrl();
|
||||
}
|
||||
|
||||
@@ -41,18 +41,18 @@ import org.apache.pulsar.client.api.PulsarClient;
|
||||
import org.apache.pulsar.client.api.Schema;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import org.springframework.pulsar.core.AbstractContainerBaseTests;
|
||||
import org.springframework.pulsar.core.DefaultPulsarConsumerFactory;
|
||||
import org.springframework.pulsar.core.DefaultPulsarProducerFactory;
|
||||
import org.springframework.pulsar.core.PulsarOperations;
|
||||
import org.springframework.pulsar.core.PulsarTemplate;
|
||||
import org.springframework.pulsar.core.PulsarTestContainerSupport;
|
||||
import org.springframework.pulsar.core.TypedMessageBuilderCustomizer;
|
||||
import org.springframework.util.backoff.FixedBackOff;
|
||||
|
||||
/**
|
||||
* @author Soby Chacko
|
||||
*/
|
||||
public class DefaultPulsarConsumerErrorHandlerTests extends AbstractContainerBaseTests {
|
||||
public class DefaultPulsarConsumerErrorHandlerTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
@@ -61,7 +61,8 @@ public class DefaultPulsarConsumerErrorHandlerTests extends AbstractContainerBas
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-1"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-1");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
|
||||
@@ -111,7 +112,8 @@ public class DefaultPulsarConsumerErrorHandlerTests extends AbstractContainerBas
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-2"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-2");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
|
||||
@@ -159,7 +161,8 @@ public class DefaultPulsarConsumerErrorHandlerTests extends AbstractContainerBas
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-3"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-3");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
|
||||
@@ -217,7 +220,8 @@ public class DefaultPulsarConsumerErrorHandlerTests extends AbstractContainerBas
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-4"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-4");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
|
||||
@@ -288,7 +292,8 @@ public class DefaultPulsarConsumerErrorHandlerTests extends AbstractContainerBas
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-5"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-5");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
|
||||
@@ -357,7 +362,8 @@ public class DefaultPulsarConsumerErrorHandlerTests extends AbstractContainerBas
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-6"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-6");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
|
||||
@@ -426,7 +432,8 @@ public class DefaultPulsarConsumerErrorHandlerTests extends AbstractContainerBas
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-7"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-7");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
|
||||
@@ -494,7 +501,8 @@ public class DefaultPulsarConsumerErrorHandlerTests extends AbstractContainerBas
|
||||
config.put("topicNames", Collections.singleton("default-error-handler-tests-8"));
|
||||
config.put("subscriptionName", "default-error-handler-tests-sub-8");
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
|
||||
|
||||
@@ -43,17 +43,17 @@ import org.apache.pulsar.client.api.SubscriptionType;
|
||||
import org.apache.pulsar.client.impl.MultiplierRedeliveryBackoff;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import org.springframework.pulsar.core.AbstractContainerBaseTests;
|
||||
import org.springframework.pulsar.core.ConsumerTestUtils;
|
||||
import org.springframework.pulsar.core.DefaultPulsarConsumerFactory;
|
||||
import org.springframework.pulsar.core.DefaultPulsarProducerFactory;
|
||||
import org.springframework.pulsar.core.PulsarTemplate;
|
||||
import org.springframework.pulsar.core.PulsarTestContainerSupport;
|
||||
|
||||
/**
|
||||
* @author Soby Chacko
|
||||
* @author Alexander Preuß
|
||||
*/
|
||||
class DefaultPulsarMessageListenerContainerTests extends AbstractContainerBaseTests {
|
||||
class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerSupport {
|
||||
|
||||
@Test
|
||||
void basicDefaultConsumer() throws Exception {
|
||||
@@ -62,7 +62,8 @@ class DefaultPulsarMessageListenerContainerTests extends AbstractContainerBaseTe
|
||||
strings.add("dpmlct-012");
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "dpmlct-sb-012");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
@@ -92,7 +93,8 @@ class DefaultPulsarMessageListenerContainerTests extends AbstractContainerBaseTe
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "dpmlct-sb-013");
|
||||
config.put("subscriptionInitialPosition", SubscriptionInitialPosition.Earliest);
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
CountDownLatch latch = new CountDownLatch(5);
|
||||
@@ -125,7 +127,8 @@ class DefaultPulsarMessageListenerContainerTests extends AbstractContainerBaseTe
|
||||
strings.add("dpmlct-014");
|
||||
config.put("topicNames", strings);
|
||||
config.put("subscriptionName", "dpmlct-sb-014");
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
|
||||
@@ -164,7 +167,8 @@ class DefaultPulsarMessageListenerContainerTests extends AbstractContainerBaseTe
|
||||
.maxDelayMs(5 * 1000).build();
|
||||
config.put("negativeAckRedeliveryBackoff", redeliveryBackoff);
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<String> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
CountDownLatch latch = new CountDownLatch(10);
|
||||
@@ -214,7 +218,8 @@ class DefaultPulsarMessageListenerContainerTests extends AbstractContainerBaseTe
|
||||
.build();
|
||||
config.put("deadLetterPolicy", deadLetterPolicy);
|
||||
|
||||
final PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(getPulsarBrokerUrl()).build();
|
||||
final PulsarClient pulsarClient = PulsarClient.builder()
|
||||
.serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build();
|
||||
final DefaultPulsarConsumerFactory<Integer> pulsarConsumerFactory = new DefaultPulsarConsumerFactory<>(
|
||||
pulsarClient, config);
|
||||
|
||||
|
||||
@@ -56,13 +56,13 @@ import org.springframework.pulsar.config.PulsarClientConfiguration;
|
||||
import org.springframework.pulsar.config.PulsarClientFactoryBean;
|
||||
import org.springframework.pulsar.config.PulsarListenerContainerFactory;
|
||||
import org.springframework.pulsar.config.PulsarListenerEndpointRegistry;
|
||||
import org.springframework.pulsar.core.AbstractContainerBaseTests;
|
||||
import org.springframework.pulsar.core.DefaultPulsarConsumerFactory;
|
||||
import org.springframework.pulsar.core.DefaultPulsarProducerFactory;
|
||||
import org.springframework.pulsar.core.PulsarAdministration;
|
||||
import org.springframework.pulsar.core.PulsarConsumerFactory;
|
||||
import org.springframework.pulsar.core.PulsarProducerFactory;
|
||||
import org.springframework.pulsar.core.PulsarTemplate;
|
||||
import org.springframework.pulsar.core.PulsarTestContainerSupport;
|
||||
import org.springframework.pulsar.core.PulsarTopic;
|
||||
import org.springframework.test.annotation.DirtiesContext;
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
@@ -75,7 +75,7 @@ import org.springframework.util.backoff.FixedBackOff;
|
||||
*/
|
||||
@SpringJUnitConfig
|
||||
@DirtiesContext
|
||||
public class PulsarListenerTests extends AbstractContainerBaseTests {
|
||||
public class PulsarListenerTests implements PulsarTestContainerSupport {
|
||||
|
||||
static CountDownLatch latch = new CountDownLatch(1);
|
||||
static CountDownLatch latch1 = new CountDownLatch(3);
|
||||
@@ -105,7 +105,7 @@ public class PulsarListenerTests extends AbstractContainerBaseTests {
|
||||
|
||||
@Bean
|
||||
public PulsarClientConfiguration pulsarClientConfiguration() {
|
||||
return new PulsarClientConfiguration(Map.of("serviceUrl", getPulsarBrokerUrl()));
|
||||
return new PulsarClientConfiguration(Map.of("serviceUrl", PulsarTestContainerSupport.getPulsarBrokerUrl()));
|
||||
}
|
||||
|
||||
@Bean
|
||||
@@ -129,7 +129,8 @@ public class PulsarListenerTests extends AbstractContainerBaseTests {
|
||||
|
||||
@Bean
|
||||
PulsarAdministration pulsarAdministration() {
|
||||
return new PulsarAdministration(PulsarAdmin.builder().serviceHttpUrl(getHttpServiceUrl()));
|
||||
return new PulsarAdministration(
|
||||
PulsarAdmin.builder().serviceHttpUrl(PulsarTestContainerSupport.getHttpServiceUrl()));
|
||||
}
|
||||
|
||||
@Bean
|
||||
|
||||
@@ -37,13 +37,13 @@ import org.springframework.pulsar.config.ConcurrentPulsarListenerContainerFactor
|
||||
import org.springframework.pulsar.config.PulsarClientConfiguration;
|
||||
import org.springframework.pulsar.config.PulsarClientFactoryBean;
|
||||
import org.springframework.pulsar.config.PulsarListenerContainerFactory;
|
||||
import org.springframework.pulsar.core.AbstractContainerBaseTests;
|
||||
import org.springframework.pulsar.core.DefaultPulsarConsumerFactory;
|
||||
import org.springframework.pulsar.core.DefaultPulsarProducerFactory;
|
||||
import org.springframework.pulsar.core.PulsarAdministration;
|
||||
import org.springframework.pulsar.core.PulsarConsumerFactory;
|
||||
import org.springframework.pulsar.core.PulsarProducerFactory;
|
||||
import org.springframework.pulsar.core.PulsarTemplate;
|
||||
import org.springframework.pulsar.core.PulsarTestContainerSupport;
|
||||
|
||||
import io.micrometer.common.KeyValues;
|
||||
import io.micrometer.core.tck.MeterRegistryAssert;
|
||||
@@ -62,7 +62,7 @@ import io.micrometer.tracing.test.simple.SpansAssert;
|
||||
* @author Chris Bono
|
||||
* @see SampleTestRunner
|
||||
*/
|
||||
public class ObservationIntegrationTests extends SampleTestRunner {
|
||||
public class ObservationIntegrationTests extends SampleTestRunner implements PulsarTestContainerSupport {
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Override
|
||||
@@ -110,7 +110,7 @@ public class ObservationIntegrationTests extends SampleTestRunner {
|
||||
|
||||
@Configuration(proxyBeanMethods = false)
|
||||
@EnablePulsar
|
||||
static class ObservationIntegrationTestAppConfig extends AbstractContainerBaseTests {
|
||||
static class ObservationIntegrationTestAppConfig {
|
||||
|
||||
@Bean
|
||||
public PulsarProducerFactory<String> pulsarProducerFactory(PulsarClient pulsarClient) {
|
||||
@@ -124,7 +124,7 @@ public class ObservationIntegrationTests extends SampleTestRunner {
|
||||
|
||||
@Bean
|
||||
public PulsarClientConfiguration pulsarClientConfiguration() {
|
||||
return new PulsarClientConfiguration(Map.of("serviceUrl", getPulsarBrokerUrl()));
|
||||
return new PulsarClientConfiguration(Map.of("serviceUrl", PulsarTestContainerSupport.getPulsarBrokerUrl()));
|
||||
}
|
||||
|
||||
@Bean
|
||||
@@ -150,7 +150,8 @@ public class ObservationIntegrationTests extends SampleTestRunner {
|
||||
|
||||
@Bean
|
||||
PulsarAdministration pulsarAdministration() {
|
||||
return new PulsarAdministration(PulsarAdmin.builder().serviceHttpUrl(getHttpServiceUrl()));
|
||||
return new PulsarAdministration(
|
||||
PulsarAdmin.builder().serviceHttpUrl(PulsarTestContainerSupport.getHttpServiceUrl()));
|
||||
}
|
||||
|
||||
@Bean
|
||||
|
||||
@@ -43,13 +43,13 @@ import org.springframework.pulsar.config.ConcurrentPulsarListenerContainerFactor
|
||||
import org.springframework.pulsar.config.PulsarClientConfiguration;
|
||||
import org.springframework.pulsar.config.PulsarClientFactoryBean;
|
||||
import org.springframework.pulsar.config.PulsarListenerContainerFactory;
|
||||
import org.springframework.pulsar.core.AbstractContainerBaseTests;
|
||||
import org.springframework.pulsar.core.DefaultPulsarConsumerFactory;
|
||||
import org.springframework.pulsar.core.DefaultPulsarProducerFactory;
|
||||
import org.springframework.pulsar.core.PulsarAdministration;
|
||||
import org.springframework.pulsar.core.PulsarConsumerFactory;
|
||||
import org.springframework.pulsar.core.PulsarProducerFactory;
|
||||
import org.springframework.pulsar.core.PulsarTemplate;
|
||||
import org.springframework.pulsar.core.PulsarTestContainerSupport;
|
||||
import org.springframework.test.annotation.DirtiesContext;
|
||||
import org.springframework.test.context.junit.jupiter.SpringJUnitConfig;
|
||||
|
||||
@@ -83,7 +83,7 @@ import io.micrometer.tracing.test.simple.SimpleTracer;
|
||||
*/
|
||||
@SpringJUnitConfig
|
||||
@DirtiesContext
|
||||
public class ObservationTests extends AbstractContainerBaseTests {
|
||||
public class ObservationTests implements PulsarTestContainerSupport {
|
||||
|
||||
private static final String LISTENER_ID_TAG = "spring.pulsar.listener.id";
|
||||
|
||||
@@ -179,7 +179,7 @@ public class ObservationTests extends AbstractContainerBaseTests {
|
||||
|
||||
@Bean
|
||||
public PulsarClientConfiguration pulsarClientConfiguration() {
|
||||
return new PulsarClientConfiguration(Map.of("serviceUrl", getPulsarBrokerUrl()));
|
||||
return new PulsarClientConfiguration(Map.of("serviceUrl", PulsarTestContainerSupport.getPulsarBrokerUrl()));
|
||||
}
|
||||
|
||||
@Bean(name = "observationTestsTemplate")
|
||||
@@ -223,7 +223,8 @@ public class ObservationTests extends AbstractContainerBaseTests {
|
||||
|
||||
@Bean
|
||||
PulsarAdministration pulsarAdministration() {
|
||||
return new PulsarAdministration(PulsarAdmin.builder().serviceHttpUrl(getHttpServiceUrl()));
|
||||
return new PulsarAdministration(
|
||||
PulsarAdmin.builder().serviceHttpUrl(PulsarTestContainerSupport.getHttpServiceUrl()));
|
||||
}
|
||||
|
||||
@Bean
|
||||
|
||||
@@ -4,7 +4,7 @@
|
||||
<pattern>%d{HH:mm:ss.SSS} [%thread] %-5level %logger - %msg%n</pattern>
|
||||
</encoder>
|
||||
</appender>
|
||||
<root level="info">
|
||||
<root level="WARN">
|
||||
<appender-ref ref="STDOUT"/>
|
||||
</root>
|
||||
<logger name="org.testcontainers" level="ERROR"/>
|
||||
|
||||
Reference in New Issue
Block a user