test: avoid real kafka clients in mq composition
This commit is contained in:
@@ -23,20 +23,29 @@ import java.util.Map;
|
||||
@ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq", name = "enabled", havingValue = "true", matchIfMissing = true)
|
||||
public class SinkMqConsumerAutoConfiguration {
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
public KafkaEnvelopeConsumerFactory kafkaEnvelopeConsumerFactory() {
|
||||
return new KafkaEnvelopeConsumerFactory();
|
||||
}
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
@ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq.consumer", name = "enabled", havingValue = "true")
|
||||
public KafkaEnvelopeConsumerRunner kafkaEnvelopeConsumerRunner(ListableBeanFactory beanFactory,
|
||||
KafkaEnvelopeConsumerFactory consumerFactory,
|
||||
SinkMqProperties props) {
|
||||
return new KafkaEnvelopeConsumerRunner(
|
||||
() -> createWorkers(beanFactory, props),
|
||||
() -> createWorkers(beanFactory, consumerFactory, props),
|
||||
Duration.ofMillis(props.getConsumer().getPollTimeoutMillis()),
|
||||
Duration.ofMillis(props.getConsumer().getLoopBackoffMillis()),
|
||||
props.getConsumer().isAutoStartup());
|
||||
}
|
||||
|
||||
private List<KafkaEnvelopeConsumerWorker> createWorkers(ListableBeanFactory beanFactory, SinkMqProperties props) {
|
||||
private List<KafkaEnvelopeConsumerWorker> createWorkers(ListableBeanFactory beanFactory,
|
||||
KafkaEnvelopeConsumerFactory consumerFactory,
|
||||
SinkMqProperties props) {
|
||||
Map<String, EnvelopeConsumerProcessor> processors = beanFactory.getBeansOfType(EnvelopeConsumerProcessor.class);
|
||||
return new KafkaEnvelopeConsumerFactory().createWorkers(processors, props);
|
||||
return consumerFactory.createWorkers(processors, props);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,20 +2,32 @@ package com.lingniu.ingest.sink.mq;
|
||||
|
||||
import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor;
|
||||
import com.lingniu.ingest.api.consumer.EnvelopeIngestResult;
|
||||
import org.apache.kafka.clients.consumer.MockConsumer;
|
||||
import org.apache.kafka.clients.consumer.OffsetResetStrategy;
|
||||
import org.apache.kafka.clients.producer.KafkaProducer;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.extension.ExtendWith;
|
||||
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
|
||||
import org.springframework.boot.test.system.CapturedOutput;
|
||||
import org.springframework.boot.test.system.OutputCaptureExtension;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.mockito.Mockito.mock;
|
||||
|
||||
@ExtendWith(OutputCaptureExtension.class)
|
||||
class SinkMqConsumerAutoConfigurationTest {
|
||||
|
||||
private final ApplicationContextRunner contextRunner = new ApplicationContextRunner()
|
||||
.withUserConfiguration(SinkMqAutoConfiguration.class, SinkMqConsumerAutoConfiguration.class)
|
||||
.withAllowBeanDefinitionOverriding(true)
|
||||
.withBean("vehicleStateEnvelopeConsumerProcessor", EnvelopeConsumerProcessor.class,
|
||||
() -> new EnvelopeConsumerProcessor(
|
||||
"vehicle-state",
|
||||
bytes -> EnvelopeIngestResult.processed("evt-1", "VIN001"),
|
||||
record -> {}))
|
||||
.withBean("kafkaProducer", KafkaProducer.class, SinkMqConsumerAutoConfigurationTest::kafkaProducer)
|
||||
.withBean(KafkaEnvelopeConsumerFactory.class,
|
||||
() -> new KafkaEnvelopeConsumerFactory(props -> new MockConsumer<>(OffsetResetStrategy.EARLIEST)))
|
||||
.withPropertyValues(
|
||||
"lingniu.ingest.sink.mq.enabled=true",
|
||||
"lingniu.ingest.sink.mq.type=kafka",
|
||||
@@ -25,10 +37,16 @@ class SinkMqConsumerAutoConfigurationTest {
|
||||
"lingniu.ingest.sink.mq.consumer.bindings.vehicleStateEnvelopeConsumerProcessor.topics[0]=vehicle.realtime");
|
||||
|
||||
@Test
|
||||
void createsKafkaEnvelopeConsumerRunnerWhenConsumerBindingIsConfigured() {
|
||||
void createsKafkaEnvelopeConsumerRunnerWhenConsumerBindingIsConfigured(CapturedOutput output) {
|
||||
contextRunner.run(context -> {
|
||||
assertThat(context).hasSingleBean(KafkaEnvelopeConsumerRunner.class);
|
||||
assertThat(context.getBean(KafkaEnvelopeConsumerRunner.class).workers()).hasSize(1);
|
||||
});
|
||||
assertThat(output).doesNotContain("114.55.58.251");
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static KafkaProducer<String, byte[]> kafkaProducer() {
|
||||
return mock(KafkaProducer.class);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user