【问题标题】:How to set port in KafkaEmbedded when unit testing spring-kafka consumer单元测试spring-kafka消费者时如何在KafkaEmbedded中设置端口
【发布时间】:2019-05-20 17:12:11
【问题描述】:

我在使用来自Kakfa 0.9 集群的消息的应用程序中使用spring-boot-starter-parent 版本1.5.0.RELEASEspring-kafka 版本1.0.0.RELEASEspring-kafka-test 版本1.0.0.RELEASE。我为我的消费者进行了单元测试,它使用了KafkaEmbedded,但由于代理端口是随机选择的,所以它失败了。有没有办法可以在不更改版本的情况下设置此代理属性?或者我应该使用哪些版本以免破坏任何东西?

这是KafkaListenerKafkaConsumerTest 的代码。

Listener.java

@Service
public class Listener {

    private static final Logger logger = LoggerFactory.getLogger(Listener.class);
    private CountDownLatch latch = new CountDownLatch(1);

    @KafkaListener(topics = "topic", group = "group", containerFactory = "kafkaListenerContainerFactory")
    public void consumeClicks(@Payload String msg, @Header(KafkaHeaders.RECEIVED_PARTITION_ID) Integer partition, @Header(KafkaHeaders.OFFSET) Integer offset, Acknowledgment ack) throws Exception {
        logger.info(msg);
        latch.countDown();
        ack.acknowledge();
    }

    public CountDownLatch getLatch() {
        return latch;
    }
}

KafkaConsumerTest.java编辑

@DirtiesContext
@SpringBootTest(classes = {SpringApplication.class})
@RunWith(SpringRunner.class)
public class KafkaConsumerTest {
    private static final Logger logger = LoggerFactory.getLogger(KafkaConsumerTest.class);
    private static String TEST_TOPIC = "topic";

    @ClassRule
    public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, TEST_TOPIC);

    public KafkaTemplate<String, String> template;

    @Autowired
    private KafkaListenerEndpointRegistry kafkaListenerEndpointRegistry;

    @Autowired
    private Listener listener;

    @Before
    public void init(){
        System.setProperty("spring.kafka.bootstrap-servers", embeddedKafka.getBrokersAsString());
        Map<String, Object> senderProps = KafkaTestUtils.senderProps(embeddedKafka.getBrokersAsString());
        senderProps.put("key.serializer", StringSerializer.class);
        ProducerFactory<String, String> producerFactory = new DefaultKafkaProducerFactory<String, String>(senderProps);
        template = new KafkaTemplate<>(producerFactory);
        template.setDefaultTopic(TEST_TOPIC);
    }

    @Test
    public void testConsume() throws Exception {
        String record = "message";
        template.sendDefault(TEST_TOPIC, record);
        logger.debug("test-consume sent record {}", record);
        listener.getLatch().await(1000, TimeUnit.MILLISECONDS);
        Assert.assertEquals(listener.getLatch().getCount(), 0);
    }
}

【问题讨论】:

    标签: java spring-kafka spring-kafka-test


    【解决方案1】:

    请使用 spring-kafka 1.3.9 和 boot 1.5;不再支持早期版本。当前启动 1.5.x 版本是 1.5.21。

    @ClassRule
    public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, TEST_TOPIC);
    
    static {
        embeddedKafka.setKafkaPorts(1234);
    }
    

    setKafkaPorts 从 1.3 开始可用。

    但是,您在测试中正确使用了分配的随机端口

    Map<String, Object> senderProps = KafkaTestUtils.senderProps(embeddedKafka.getBrokersAsString());
    

    要让 kafka 监听器连接到嵌入式代理,您可以使用。

        System.setProperty("spring.kafka.bootstrap-servers", embeddedKafka.getBrokersAsString()); 
    

    【讨论】:

    • 我的应用程序应该可以与Kafka 0.9 一起使用,而且我认为这种情况不会很快改变。您提到的版本可以使用吗?
    • 0.9 太旧了;有很多改进,尤其是KIP-62,它解决了慢客户端的再平衡问题,你至少应该尝试迁移到 0.10.2.0。当前版本是 2.2.0。 compatibility page 甚至没有提到 0.9.0.0。我不知道 1.3.9 是否适用于 0.9.0.0;你应该试试看。
    • 迁移到 0.10.2.0 正在与我公司的 DevOps 团队进行谈判,不幸的是,这将是一个缓慢的过程。我将我的应用程序与 1.3.9 一起使用,并运行本地 kafka 0.9.0.0 集群。应用程序只打印消费者配置值和o.a.k.c.u.AppInfoParser : Kafka version : 0.11.0.2 并停止。不进行分区分配。我想我现在必须使用旧版本。
    • 您至少应该升级到 1.0.6.RELEASE,这是 1.0.x 的最后一个版本,用于修复错误。
    • 使用 1.0.6,我将System.setProperty(..) 添加到测试类的init() 方法中,但KafkaListener 仍然连接到在consumerConfig 中为侦听器容器配置的代理,测试失败。
    【解决方案2】:

    我认为在为测试加载应用程序上下文时,正在创建两个类型为 (ProducerFactory and KafkaTemplate) 的 bean,一个是原始配置,第二个是测试配置,尝试使用不同的配置文件进行测试 application-test.yml 并添加 bean 覆盖属性

    spring.main.allow-bean-definition-overriding to true.
    

    这样它将使用测试 bean 覆盖应用程序 bean,并将 ProducerFactoryKafkaTemplate 声明为测试中的 bean,与应用程序中的名称相同

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-05-09
      • 1970-01-01
      • 1970-01-01
      • 2019-04-22
      • 1970-01-01
      • 2018-10-20
      • 1970-01-01
      • 2021-05-20
      相关资源
      最近更新 更多