From 6ab2f2538738048c0cc7b06959f106766e22fdc1 Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 03:54:14 +0800 Subject: [PATCH] test: avoid real kafka clients in mq composition --- .../mq/SinkMqConsumerAutoConfiguration.java | 15 +++++++++++--- .../SinkMqConsumerAutoConfigurationTest.java | 20 ++++++++++++++++++- 2 files changed, 31 insertions(+), 4 deletions(-) diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqConsumerAutoConfiguration.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqConsumerAutoConfiguration.java index 2302d080..025fc10c 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqConsumerAutoConfiguration.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqConsumerAutoConfiguration.java @@ -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 createWorkers(ListableBeanFactory beanFactory, SinkMqProperties props) { + private List createWorkers(ListableBeanFactory beanFactory, + KafkaEnvelopeConsumerFactory consumerFactory, + SinkMqProperties props) { Map processors = beanFactory.getBeansOfType(EnvelopeConsumerProcessor.class); - return new KafkaEnvelopeConsumerFactory().createWorkers(processors, props); + return consumerFactory.createWorkers(processors, props); } } diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/SinkMqConsumerAutoConfigurationTest.java b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/SinkMqConsumerAutoConfigurationTest.java index 88b196a0..72126b59 100644 --- a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/SinkMqConsumerAutoConfigurationTest.java +++ b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/SinkMqConsumerAutoConfigurationTest.java @@ -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 kafkaProducer() { + return mock(KafkaProducer.class); } }