【问题标题】:How do I Restart a shutdown embeddedKafkaServer in a Spring Unit Test?如何在 Spring 单元测试中重新启动关闭的嵌入式KafkaServer?
【发布时间】:2020-09-30 20:59:16
【问题描述】:

我有一个 Spring-boot 单元测试,用于在主 Kafka 集群上线时测试我的应用程序的切换回功能。

当主节点离线时,应用程序成功切换到辅助节点。现在我们添加了在计时器而不是故障时切换回主节点的功能。

我的测试方法如下:

   //Rochelle = Primary BootStrapServers
   //Hudson   = Secondary BootStrapServers


   @Test
   public void send_switchback() throws Exception
   {
      //Get ABSwitchCluster to check failover details
      KafkaSwitchCluster ktSwitch = (KafkaSwitchCluster)
              ((BootStrapExposerProducerFactory)
                       kafkaTemplate.getProducerFactory()).getBootStrapSupplier();

      assertThat(ktSwitch,             notNullValue());
      assertThat(ktSwitch.get(),       is(Rochelle));
      assertThat(ktSwitch.isPrimary(), is(true));

      assertThat(getBootStrapServersList(), is(Rochelle));

      log.info("Shutdown Broker to test Failover.");

      //Shutdown Primary Servers to simulate disconnection
      shutdownBroker_primary();
      //Allow for fail over to happen
      if ( ktSwitch.isPrimary() )
      {
         try
         {
            synchronized (lock)
            {  //pause to give Idle Event a chance to fire
               for (int i = 0; i <= timeOut && ktSwitch.isPrimary(); ++i)
               //while ( ktSwitch.isPrimary() )
               {  //poll for cluster switch
                  lock.wait(Duration.ofSeconds(15).toMillis());
               }
            }
         }
         catch (InterruptedException IGNORE)
         { fail("Unable to wait for cluster switch. " + IGNORE.getMessage()); }
      }

      //Confirm Failover has happened
      assertThat(ktSwitch.get(),            is(Hudson));
      assertThat(ktSwitch.isPrimary(),      is(false));
      assertThat(getBootStrapServersList(), is(Hudson));

      assertThat(kafkaSwitchCluster.get(),       is(Hudson));
      assertThat(kafkaSwitchCluster.isPrimary(), is(false));

      //Send a message on backup server
      String message = "Test Failover";
      send(message);

      String msg = records.poll(10, TimeUnit.SECONDS);
      assertThat(msg, notNullValue());
      assertThat(msg, is(message));

      startup_primary();
      //embeddedKafkaRule.getEmbeddedKafka();

      assertThat(embeddedKafka.getBrokersAsString(), is(Rochelle));
      String brokers = embeddedKafka.getBrokersAsString();

      if ( !kafkaProducerErrorHandler.areBrokersUp(brokers) )
      {
         synchronized (lock)
         {
            for ( int i=0;
                  i <= 15 && !kafkaProducerErrorHandler.areBrokersUp(brokers)
                  && registry.isRunning();
                  ++i )
            { lock.wait(Duration.ofSeconds(1).toMillis()); }
         }
      }

      //TODO: test Scheduled Fire
      kafkaProducerErrorHandler.primarySwitch();

      if ( !kafkaSwitchCluster.isPrimary() )
      {
         try
         {
            synchronized (lock)
            {  //pause to give Idle Event a chance to fire
               for (int i = 0; i <= timeOut && !kafkaSwitchCluster.isPrimary(); ++i)
               //while ( !ktSwitch.isPrimary() )
               {  //poll for cluster switch
                  lock.wait(Duration.ofSeconds(15).toMillis());
               }
            }
         }
         catch (InterruptedException IGNORE)
         { fail("Unable to wait for cluster switch. " + IGNORE.getMessage()); }
      }

      assertThat(brokers,              anyOf(is(Rochelle), is(Hudson))); //port didn't change
      assertThat(brokers,              is(Rochelle)); //is primary
      assertThat(kafkaSwitchCluster.isPrimary(), is(true));
      //assertThat(ktSwitch.isPrimary(), is(true));
      assertThat(ktSwitch.get(),       is(brokers));

      assertThat(kafkaProducerErrorHandler.areBrokersUp(brokers),  is(true));
      assertThat(kafkaProducerErrorHandler.areBrokersUp(Rochelle), is(true));

      assertThat(ktSwitch.isPrimary(), is(true));
      //assertThat(ktSwitch.get(),       not(anyOf(is(Hudson), is(Rochelle))));
      assertThat(ktSwitch.get(),       is(embeddedKafka.getBrokersAsString()));

      //Send a message on backup server
      message = "Test newPrimary";
      send(message);

      msg = records.poll(10, TimeUnit.SECONDS);
      assertThat(msg, notNullValue());
      assertThat(msg, is(message));

      log.info("Test is finished");
   }

我正在使用这种方法来关闭我的 Primary Embedded Kafka

   public void shutdownBroker_primary()
   {
      for(KafkaServer ks : embeddedKafka.getKafkaServers())
      { ks.shutdown(); }
      for(KafkaServer ks : embeddedKafka.getKafkaServers())
      { ks.awaitShutdown(); }
   }

我正在使用它来重新启动 Kafka:

public void startup_primary()
   {
      //registry.stop();
      //kafkaSwitchCluster.Rochelle = embeddedKafka.getBrokersAsString();
      for(KafkaServer ks : embeddedKafka.getKafkaServers()) { ks.startup(); }
      registry.start();
   }

primarySwitch() 是一个计划事件,用于将集群切换回主集群。在测试中直接调用。它是在 Kafka 宕机时切换正在使用的集群的相同代码的包装器。

如何在关闭主嵌入式 Kafka 集群后成功启动它,以便我可以证明应用程序可以在主集群再次可用时成功移回主集群?


更新:
我已经在 Github 上创建了我目前拥有的代码示例:https://github.com/raystorm/Kafka-Example


更新:2Linked Repository Above 已根据下面接受的答案进行了更新,现在所有测试都通过了。

【问题讨论】:

    标签: java spring spring-boot spring-kafka


    【解决方案1】:

    它并不是真正为这个用例设计的,但只要您不需要在代理实例之间保留数据,以下方法就可以工作......

    @SpringBootTest
    @EmbeddedKafka(topics = "so64145670", bootstrapServersProperty = "spring.kafka.bootstrap-servers")
    class So64145670ApplicationTests {
    
        @Autowired
        private EmbeddedKafkaBroker broker;
    
        @Test
        void restartBroker(@Autowired KafkaTemplate<String, String> template) throws Exception {
            SendResult<String, String> sendResult = template.send("so64145670", "foo").get(10, TimeUnit.SECONDS);
            System.out.println("+++" + sendResult.getRecordMetadata());
            this.broker.destroy();
            // restart
            this.broker.afterPropertiesSet();
            sendResult = template.send("so64145670", "bar").get(10, TimeUnit.SECONDS);
            System.out.println("+++" + sendResult.getRecordMetadata());
        }
    
    }
    

    编辑

    这是一个有两个经纪人的...

    @SpringBootTest(classes = { So64145670Application.class, So64145670ApplicationTests.Config.class })
    @EmbeddedKafka(topics = "so64145670", bootstrapServersProperty = "spring.kafka.bootstrap-servers")
    class So64145670ApplicationTests {
    
        @Autowired
        private EmbeddedKafkaBroker embeddedKafka;
    
        @Autowired
        private EmbeddedKafkaBroker secondBroker;
    
        @Test
        void restartBroker(@Autowired KafkaTemplate<String, String> template,
                @Autowired ProducerFactory<String, String> pf) throws Exception {
    
            SendResult<String, String> sendResult = template.send("so64145670", "foo").get(10, TimeUnit.SECONDS);
            System.out.println("+++" + sendResult.getRecordMetadata());
            KafkaTemplate<String, String> secondTemplate = new KafkaTemplate<>(pf,
                    Map.of(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, this.secondBroker.getBrokersAsString()));
            sendResult = secondTemplate.send("so64145670-1", "foo").get(10, TimeUnit.SECONDS);
            System.out.println("+++" + sendResult.getRecordMetadata());
            this.embeddedKafka.destroy();
            this.secondBroker.destroy();
            // restart
            this.embeddedKafka.afterPropertiesSet();
            this.secondBroker.afterPropertiesSet();
            sendResult = template.send("so64145670", "bar").get(10, TimeUnit.SECONDS);
            System.out.println("+++" + sendResult.getRecordMetadata());
            sendResult = secondTemplate.send("so64145670-1", "bar").get(10, TimeUnit.SECONDS);
            System.out.println("+++" + sendResult.getRecordMetadata());
        }
    
        @Configuration
        public static class Config {
    
            @Bean
            EmbeddedKafkaBroker secondBroker() {
                return new EmbeddedKafkaBroker(1, true, "so64145670-1")
                        .brokerListProperty("spring.kafka.second.server");
            }
    
        }
    
    }
    
    +++so64145670-1@0
    +++so64145670-1-0@0
    +++so64145670-1@0
    +++so64145670-1-0@0
    

    【讨论】:

    • afterPropertiesSet() 为我抛出异常。 KafkaException: Socket server failed to bind to localhost:50674: Address already in use.
    • 有什么不对,destroy会关闭旧服务器,等待完成;如果您仍然手动调用shutdown(),您还需要调用awaitShutdown(并关闭嵌入式zookeeper服务器) - 调用destroy()更容易。
    • 删除了 shutdown()awaitShutdown() 以支持 destroy()。同样的例外。但我的测试是运行 2 个嵌入式代理。 embeddedKafkaembeddedKafka_secondary 他们会不会互相干扰?
    • 没有什么我能想到的;如果你能提供一个小的、完整的、精简的例子,我可以看看。
    • 您有 3 个代理 - 2 个 @ClassRules 和一个 @EmbeddedKafka - 将 bean 添加到测试上下文。您应该使用 Spring bean 或类规则;不是两者都 - 删除@EmbeddedKafka 或简单地将第二个代理添加为@Bean。目前尚不完全清楚为什么会导致地址已在使用中;你最终应该有 3 个不同的端口,但是让我们看看当你清理它时会发生什么。此外,以后在发布示例时,请避免使用 Lombok - 对于我们这些不使用 Lombok 开始查看示例正在做什么的人来说,这是一种痛苦。
    猜你喜欢
    • 1970-01-01
    • 2010-11-06
    • 2020-07-15
    • 1970-01-01
    • 1970-01-01
    • 2023-03-10
    • 2020-12-13
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多