【问题标题】:TimeoutException when producing event from callback从回调产生事件时出现 TimeoutException
【发布时间】:2021-04-13 12:15:27
【问题描述】:

当我尝试从另一个发送的ListenableFutureCallback 中发送消息时,我遇到了生产者超时,但如果我同步等待结果并且不使用回调,则不会发生这种情况。几条消息后,可以看到类似的错误:

...[kafka-producer-network-thread | producer-1] ERROR o.s.k.s.LoggingProducerListener - Exception thrown when sending a message with key='null' and payload='...' to topic ...: org.apache.kafka.common.errors.TimeoutException: Topic ... not present in metadata after 60000 ms.

示例代码如下:

// This produces the TimeoutException
kafkaTemplate.send(producerRecord)
  .addCallback(new ListenableFutureCallback<SendResult<String, String>>() {
      @Override
      public void onFailure(Throwable ex) {
        kafkaTemplate.send(failureMessage);
      }

      @Override
      public void onSuccess(SendResult<String, String> result) {
        kafkaTemplate.send(successMessage); // <- timeout here
      }
    });

// This does work as expected
kafkaTemplate.send(producerRecord).get();
kafkaTemplate.send(successMessage);

我在文档中没有发现禁止从另一个生产者的回调中生成消息,或者它与其他一些配置有关?

【问题讨论】:

    标签: java apache-kafka kafka-producer-api


    【解决方案1】:

    使用下面的代码进行测试。可以在onSuccess 中成功发送消息而不会超时。你能粘贴整个示例代码吗?

    public class TestProducer {
    
        private static final Logger LOG = LoggerFactory.getLogger(TestProducer.class);
    
        // kafka-topics --list --zookeeper localhost:2181
        // kafka-topics --create --zookeeper 127.0.0.1:2181 --topic test_callback --replication-factor 1 --partitions 1
        public static void main(String[] args) throws Exception {
            Map<String, Object> configProps = new HashMap<>();
            configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
            configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
            configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
            DefaultKafkaProducerFactory<String, String> producerFactory = new DefaultKafkaProducerFactory<>(configProps);
    
            KafkaTemplate<String, String> kafkaTemplate = new KafkaTemplate<>(producerFactory);
            ProducerRecord<String, String> producerRecord = new ProducerRecord<>("test_callback", 0, "key", "data");
    
            ProducerRecord<String, String> failureMessage = new ProducerRecord<>("test_callback", 0, "failed", "failed data");
            ProducerRecord<String, String> successMessage = new ProducerRecord<>("test_callback", 0, "success", "success data");
    
            kafkaTemplate.send(producerRecord)
                    .addCallback(new ListenableFutureCallback<SendResult<String, String>>() {
                        @Override
                        public void onFailure(Throwable ex) {
                            kafkaTemplate.send(failureMessage);
                            // System.out.println("onFailure");
                        }
    
                        @Override
                        public void onSuccess(SendResult<String, String> result) {
                            kafkaTemplate.send(successMessage); // <- timeout here
                            System.out.println("onSuccess");
                        }
                    });
    
    
            // kafkaTemplate.send(producerRecord).get();
            // kafkaTemplate.send(successMessage).get();
    
            ThreadUtils.sleepQuietly(1000000);
        }
    }
    

    【讨论】:

      猜你喜欢
      • 2017-05-27
      • 1970-01-01
      • 1970-01-01
      • 2010-10-28
      • 2011-01-26
      • 1970-01-01
      • 1970-01-01
      • 2018-03-20
      • 1970-01-01
      相关资源
      最近更新 更多