PulsarReader integration tests

This commit is contained in:
Soby Chacko
2023-02-17 19:50:34 -05:00
parent b82128c82c
commit 23f99593c3
8 changed files with 224 additions and 3 deletions

View File

@@ -22,6 +22,7 @@ import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.common.schema.SchemaType;
import org.springframework.messaging.handler.annotation.MessageMapping;
@@ -38,7 +39,7 @@ public @interface PulsarReader {
* <p>
* If none is specified an auto-generated id is used.
* <p>
* SpEL {@code #{...}} and property place holders {@code ${...}} are supported.
* SpEL {@code #{...}} and property placeholders {@code ${...}} are supported.
* @return the {@code id} for the container managing for this endpoint.
* @see PulsarReaderEndpointRegistry#getReaderContainer(String)
*/
@@ -56,6 +57,12 @@ public @interface PulsarReader {
*/
SchemaType schemaType() default SchemaType.NONE;
/**
* {@link MessageId} for this reader to start from.
* @return starting message id - earliest or latest.
*/
String startMessageId() default "";
/**
* Topics to listen to.
* @return a comma separated list of topics to listen from.

View File

@@ -26,6 +26,8 @@ import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;
import org.apache.pulsar.client.api.MessageId;
import org.springframework.aop.support.AopUtils;
import org.springframework.beans.BeansException;
import org.springframework.beans.factory.BeanInitializationException;
@@ -210,6 +212,15 @@ public class PulsarReaderAnnotationBeanPostProcessor<V> extends AbstractPulsarAn
endpoint.setId(getEndpointId(pulsarReader));
endpoint.setTopics(topics);
endpoint.setSchemaType(pulsarReader.schemaType());
String startMessageIdString = pulsarReader.startMessageId();
MessageId startMessageId = null;
if (startMessageIdString.equalsIgnoreCase("earliest")) {
startMessageId = MessageId.earliest;
}
else if (startMessageIdString.equalsIgnoreCase("latest")) {
startMessageId = MessageId.latest;
}
endpoint.setStartMessageId(startMessageId);
String autoStartup = pulsarReader.autoStartup();
if (StringUtils.hasText(autoStartup)) {

View File

@@ -21,6 +21,7 @@ import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.common.schema.SchemaType;
import org.springframework.beans.BeansException;
@@ -55,6 +56,8 @@ public abstract class AbstractPulsarReaderEndpoint<K>
private final List<String> topics = new ArrayList<>();
private MessageId startMessageId;
private BeanFactory beanFactory;
private BeanExpressionResolver resolver;
@@ -170,4 +173,12 @@ public abstract class AbstractPulsarReaderEndpoint<K>
this.schemaType = schemaType;
}
public MessageId getStartMessageId() {
return this.startMessageId;
}
public void setStartMessageId(MessageId startMessageId) {
this.startMessageId = startMessageId;
}
}

View File

@@ -53,6 +53,7 @@ public class DefaultPulsarReaderContainerFactory<T>
}
properties.setSchemaType(endpoint.getSchemaType());
properties.setStartMessageId(endpoint.getStartMessageId());
return new DefaultPulsarReaderListenerContainer<>(this.getReaderFactory(), properties);
}

View File

@@ -18,6 +18,7 @@ package org.springframework.pulsar.config;
import java.util.List;
import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.common.schema.SchemaType;
import org.springframework.lang.Nullable;
@@ -78,4 +79,6 @@ public interface PulsarReaderEndpoint<C extends PulsarReaderListenerContainer> {
@Nullable
Boolean getAutoStartup();
MessageId getStartMessageId();
}

View File

@@ -150,7 +150,7 @@ public class DefaultPulsarReaderListenerContainer<T> extends AbstractPulsarReade
readerContainerProperties.getStartMessageId(), (Schema) readerContainerProperties.getSchema());
}
catch (PulsarClientException e) {
DefaultPulsarReaderListenerContainer.this.logger.error(e, () -> "Pulsar client exceptions.");
throw new IllegalStateException("Pulsar client exceptions.", e);
}
}

View File

@@ -25,6 +25,7 @@ import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.common.schema.SchemaType;
import org.springframework.core.task.AsyncTaskExecutor;
import org.springframework.pulsar.core.DefaultSchemaResolver;
import org.springframework.pulsar.core.SchemaResolver;
import org.springframework.util.Assert;
@@ -45,7 +46,7 @@ public class PulsarReaderContainerProperties {
private List<String> topics;
private MessageId startMessageId = MessageId.earliest;
private MessageId startMessageId;
private Schema<?> schema;
@@ -59,6 +60,10 @@ public class PulsarReaderContainerProperties {
return this.readerListener;
}
public PulsarReaderContainerProperties() {
this.schemaResolver = new DefaultSchemaResolver();
}
public void setReaderListener(Object readerListener) {
this.readerListener = readerListener;
}

View File

@@ -0,0 +1,183 @@
/*
* Copyright 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.pulsar.reader;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import org.apache.pulsar.client.api.PulsarClient;
import org.junit.jupiter.api.Nested;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.pulsar.annotation.EnablePulsar;
import org.springframework.pulsar.annotation.PulsarReader;
import org.springframework.pulsar.config.DefaultPulsarReaderContainerFactory;
import org.springframework.pulsar.config.PulsarClientFactoryBean;
import org.springframework.pulsar.config.PulsarReaderContainerFactory;
import org.springframework.pulsar.core.DefaultPulsarProducerFactory;
import org.springframework.pulsar.core.DefaultPulsarReaderFactory;
import org.springframework.pulsar.core.PulsarProducerFactory;
import org.springframework.pulsar.core.PulsarReaderFactory;
import org.springframework.pulsar.core.PulsarTemplate;
import org.springframework.pulsar.test.support.PulsarTestContainerSupport;
import org.springframework.test.annotation.DirtiesContext;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit.jupiter.SpringJUnitConfig;
/**
* {@link PulsarReader} integration tests.
*
* @author Soby Chacko
*/
@SpringJUnitConfig
@DirtiesContext
public class PulsarReaderTests implements PulsarTestContainerSupport {
@Autowired
PulsarTemplate<String> pulsarTemplate;
@Autowired
private PulsarClient pulsarClient;
@Configuration(proxyBeanMethods = false)
@EnablePulsar
public static class TopLevelConfig {
@Bean
public PulsarProducerFactory<String> pulsarProducerFactory(PulsarClient pulsarClient) {
Map<String, Object> config = Collections.emptyMap();
return new DefaultPulsarProducerFactory<>(pulsarClient, config);
}
@Bean
public PulsarClientFactoryBean pulsarClientFactoryBean() {
return new PulsarClientFactoryBean(Map.of("serviceUrl", PulsarTestContainerSupport.getPulsarBrokerUrl()));
}
@Bean
public PulsarTemplate<String> pulsarTemplate(PulsarProducerFactory<String> pulsarProducerFactory) {
return new PulsarTemplate<>(pulsarProducerFactory);
}
@Bean
public PulsarReaderFactory<?> pulsarReaderFactory(PulsarClient pulsarClient) {
Map<String, Object> config = new HashMap<>();
return new DefaultPulsarReaderFactory<>(pulsarClient, config);
}
@Bean
PulsarReaderContainerFactory pulsarReaderContainerFactory(PulsarReaderFactory<Object> pulsarReaderFactory) {
DefaultPulsarReaderContainerFactory<?> pulsarReaderContainerFactory = new DefaultPulsarReaderContainerFactory<>(
pulsarReaderFactory, new PulsarReaderContainerProperties());
return pulsarReaderContainerFactory;
}
}
@Nested
@ContextConfiguration(classes = StartMessageIdEarliest.PulsarReaderStartMessageIdEarliest.class)
class StartMessageIdEarliest {
static CountDownLatch latch = new CountDownLatch(1);
@Test
void startMessageIdEarliest() throws Exception {
pulsarTemplate.send("pulsarReaderBasicScenario-topic-1", "hello foo");
assertThat(latch.await(5, TimeUnit.SECONDS)).isTrue();
}
@EnablePulsar
@Configuration
static class PulsarReaderStartMessageIdEarliest {
@PulsarReader(id = "pulsarReaderBasicScenario-id-1",
subscriptionName = "pulsarReaderBasicScenario-subscription-1",
topics = "pulsarReaderBasicScenario-topic-1", startMessageId = "earliest")
void read(String ignored) {
latch.countDown();
}
}
}
@Nested
class StartMessageIdMissing {
@Test
void startMessageIdMissing() {
assertThatThrownBy(() -> new AnnotationConfigApplicationContext(TopLevelConfig.class,
PulsarReaderStartMessageIdMissing.class)).rootCause().isInstanceOf(IllegalArgumentException.class)
.hasMessage(
"Start message id or start message from roll back must be specified but they cannot be specified at the same time");
}
@EnablePulsar
static class PulsarReaderStartMessageIdMissing {
@PulsarReader(id = "pulsarReaderBasicScenario-id-2",
subscriptionName = "pulsarReaderBasicScenario-subscription-2",
topics = "pulsarReaderBasicScenario-topic-2")
void readWithoutStartMessageId(String ignored) {
}
}
}
@Nested
class StartMessageIdLatest {
static CountDownLatch latch = new CountDownLatch(1);
@Test
void startMessageIdLatest() throws Exception {
pulsarTemplate.send("pulsarReaderBasicScenario-topic-3", "hello foo");
new AnnotationConfigApplicationContext(TopLevelConfig.class,
StartMessageIdLatest.PulsarReaderStartMessageIdLatest.class);
pulsarTemplate.send("pulsarReaderBasicScenario-topic-3", "hello foobar");
assertThat(latch.await(5, TimeUnit.SECONDS)).isTrue();
}
@EnablePulsar
@Configuration
static class PulsarReaderStartMessageIdLatest {
@PulsarReader(id = "pulsarReaderBasicScenario-id-3",
subscriptionName = "pulsarReaderBasicScenario-subscription-3",
topics = "pulsarReaderBasicScenario-topic-3", startMessageId = "latest")
void read(String msg) {
latch.countDown();
assertThat(msg).isEqualTo("hello foobar");
}
}
}
}