【问题标题】:Spring Kafka Unit Tests Triggers the listener, but the method cannot get the message using consumer.pollSpring Kafka Unit Tests 触发监听,但是方法使用consumer.poll无法获取消息
【发布时间】:2020-02-19 00:21:12
【问题描述】:

我们正在使用 spring-kafka-test-2.2.8-RELEASE。 当我使用模板发送消息时,它正确触发了监听器,但在consumer.poll中无法获取消息内容。如果我实例化 KafkaTemplate 而不在类属性中“连接”它并基于生产者工厂实例化它,它会发送消息,但不会触发 @KafkaListener,只有当我在 @Test 方法中设置消息侦听器时才有效。我需要触发kafka监听器,并意识到接下来会调用哪个Topic(“成功”主题执行时没有错误,“errorTopic”监听器抛出异常)和消息内容。

    @RunWith(SpringRunner.class)
    @SpringBootTest
    @EmbeddedKafka(partitions = 1, topics = { "tp-in-gco-mao-notasfiscais" })
    public class InvoicingServiceTest {

         @Autowired
         private NFKafkaListener nfKafkaListener;

         @ClassRule
         public static EmbeddedKafkaRule broker = new EmbeddedKafkaRule(1, false, "tp-in-gco-mao- 
         notasfiscais");

         @Value("${" + EmbeddedKafkaBroker.SPRING_EMBEDDED_KAFKA_BROKERS + "}")
         private String brokerAddresses;

         @Autowired
         private KafkaTemplate<Object, Object> template;

         @BeforeClass
         public static void setup() {
                System.setProperty(EmbeddedKafkaBroker.BROKER_LIST_PROPERTY,
              "spring.kafka.bootstrap-servers");
         }

         @Test
         public void testTemplate() throws Exception {
               NFServiceTest nfServiceTest = spy(new NFServiceTest());

               nfKafkaListener.setNfServiceClient(nfServiceTest);
               Map<String, Object> consumerProps =  KafkaTestUtils.consumerProps("teste9", "false", broker.getEmbeddedKafka());
               consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
               consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
               consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, InvoiceDeserializer.class);
               consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");

               DefaultKafkaConsumerFactory<Integer, Object> cf = new DefaultKafkaConsumerFactory<Integer, Object>(
            consumerProps);

           Consumer<Integer, Object> consumer = cf.createConsumer();

           broker.getEmbeddedKafka().consumeFromAnEmbeddedTopic(consumer, "tp-in-gco-mao-notasfiscais");

           ZfifNfMao zf = new ZfifNfMao();
           zf.setItItensnf(new Zfietb011());

           Zfietb011 zfietb011 = new Zfietb011();
           Zfie011 zfie011 = new Zfie011();
           zfie011.setMatkl("TESTE");
           zfietb011.getItem().add(zfie011);
           zf.setItItensnf(zfietb011);

           template.send("tp-in-gco-mao-notasfiscais", zf);
           List<ConsumerRecord<Integer, Object>> received = new ArrayList<>();
           int n = 0;
           while (received.size() < 1 && n++ < 10) {
                ConsumerRecords<Integer, Object> records1 = consumer.poll(Duration.ofSeconds(10));
                //records1  is always empty
                if (!records1.isEmpty()) {
                    records1.forEach(rec -> received.add(rec));
                }
           }

           assertThat(received).extracting(rec -> {
               ZfifNfMao zfifNfMaoRdesponse = (ZfifNfMao) rec.value();
               return zfifNfMaoRdesponse.getItItensnf().getItem().get(0).getMatkl();
            }).contains("TESTE");
            broker.getEmbeddedKafka().getKafkaServers().forEach(b -> b.shutdown());
            broker.getEmbeddedKafka().getKafkaServers().forEach(b -> b.awaitShutdown());
            consumer.close();
        }

        public static class NFServiceTest implements INFServiceClient {
            CountDownLatch latch = new CountDownLatch(1);

            @Override
            public ZfifNfMaoResponse enviarSap(ZfifNfMao zfifNfMao) {
                ZfifNfMaoResponse zfifNfMaoResponse = new ZfifNfMaoResponse();
                zfifNfMaoResponse.setItItensnf(new Zfietb011());

                Zfietb011 zfietb011 = new Zfietb011();
                Zfie011 zfie011 = new Zfie011();
                zfie011.setMatkl("TESTE");
                zfietb011.getItem().add(zfie011);
                zfifNfMaoResponse.setItItensnf(zfietb011);
                return zfifNfMaoResponse;
            }
        }
    }     

【问题讨论】:

    标签: spring-kafka spring-kafka-test


    【解决方案1】:

    您有两个经纪人;一个由@EmbeddedKafka 创建,一个由@ClassRule 创建。

    使用其中一种;最好是 @EmbeddedKafka 和简单的 @Autowired 代理实例。

    我猜消费者正在听不同的经纪人;您可以通过查看消费者配置输出的 INFO 日志来确认这一点。

    【讨论】:

    • 我听从了你的建议,但它一直触发监听器,但 consumer.poll 没有捕获主题内容。
    • 打开org.springframework.kafka 的调试日志记录。如果您无法从中弄清楚,请将日志发布到某个地方(例如 GitHub Gist、pastebin 或类似的地方。
    • 您的生产者和消费者未连接到嵌入式代理。 bootstrap.servers = [RVMTDV1147.riachuelo.net:9092, RVMTDV1148.riachuelo.net:9092, RVMTDV1154.riachuelo.net:9092]。对于 2.3 及更高版本,您可以将 bootstrapServersProperty="spring.kafka.bootstrap-servers" 添加到 @EmbeddedKafka。对于早期版本,将@TestPropertySource(properties = "spring.kafka.bootstrap-servers = ${spring.embedded.kafka.brokers}") 添加到测试用例中。
    【解决方案2】:

    我听从了你的建议,但它一直触发监听器,但 consumer.poll 没有捕获主题内容。

    @RunWith(SpringRunner.class)
    @SpringBootTest
    @EmbeddedKafka(partitions = 1, topics = { "tp-in-gco-mao-notasfiscais" })
    public class InvoicingServiceTest {
    
         @Autowired
         private NFKafkaListener nfKafkaListener;
    
     @Autowired
     public EmbeddedKafkaBroker broker;
    
     @Autowired
     private KafkaTemplate<Object, Object> template;
    
     @BeforeClass
     public static void setup() {
            System.setProperty(EmbeddedKafkaBroker.BROKER_LIST_PROPERTY,
          "spring.kafka.bootstrap-servers");
     }
    
     @Test
     public void testTemplate() throws Exception {
           NFServiceTest nfServiceTest = spy(new NFServiceTest());
    
           nfKafkaListener.setNfServiceClient(nfServiceTest);
           Map<String, Object> consumerProps =  KafkaTestUtils.consumerProps("teste9", "false", broker);
           consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
           consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
           consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, InvoiceDeserializer.class);
    
           DefaultKafkaConsumerFactory<Integer, Object> cf = new DefaultKafkaConsumerFactory<Integer, Object>(
        consumerProps);
    
       Consumer<Integer, Object> consumer = cf.createConsumer();
    
       broker.consumeFromAnEmbeddedTopic(consumer, "tp-in-gco-mao-notasfiscais");
    
       ZfifNfMao zf = new ZfifNfMao();
       zf.setItItensnf(new Zfietb011());
    
       Zfietb011 zfietb011 = new Zfietb011();
       Zfie011 zfie011 = new Zfie011();
       zfie011.setMatkl("TESTE");
       zfietb011.getItem().add(zfie011);
       zf.setItItensnf(zfietb011);
    
       template.send("tp-in-gco-mao-notasfiscais", zf);
       List<ConsumerRecord<Integer, Object>> received = new ArrayList<>();
       int n = 0;
       while (received.size() < 1 && n++ < 10) {
            ConsumerRecords<Integer, Object> records1 = consumer.poll(Duration.ofSeconds(10));
            //records1  is always empty
            if (!records1.isEmpty()) {
                records1.forEach(rec -> received.add(rec));
            }
       }
    
       assertThat(received).extracting(rec -> {
           ZfifNfMao zfifNfMaoRdesponse = (ZfifNfMao) rec.value();
           return zfifNfMaoRdesponse.getItItensnf().getItem().get(0).getMatkl();
        }).contains("TESTE");
        broker.getKafkaServers().forEach(b -> b.shutdown());
        broker.getKafkaServers().forEach(b -> b.awaitShutdown());
        consumer.close();
    }
    
    public static class NFServiceTest implements INFServiceClient {
        CountDownLatch latch = new CountDownLatch(1);
    
        @Override
        public ZfifNfMaoResponse enviarSap(ZfifNfMao zfifNfMao) {
            ZfifNfMaoResponse zfifNfMaoResponse = new ZfifNfMaoResponse();
            zfifNfMaoResponse.setItItensnf(new Zfietb011());
    
            Zfietb011 zfietb011 = new Zfietb011();
            Zfie011 zfie011 = new Zfie011();
            zfie011.setMatkl("TESTE");
            zfietb011.getItem().add(zfie011);
            zfifNfMaoResponse.setItItensnf(zfietb011);
            return zfifNfMaoResponse;
        }
    }
    } 
    

    【讨论】:

    • 不要添加另一个答案 - 它不会回答问题;编辑问题以显示新信息。
    • 我在学习。我以为我无法删除/编辑。我什至不知道如何发布代码......下次我不会再犯这个错误了。顺便说一句,我已经在github上发布了另一个问题的日志“github.com/lucaslopescunha/logspringkafka/blob/master/log.log
    猜你喜欢
    • 1970-01-01
    • 2020-01-17
    • 1970-01-01
    • 2017-10-27
    • 2018-07-17
    • 1970-01-01
    • 1970-01-01
    • 2019-01-10
    • 1970-01-01
    相关资源
    最近更新 更多