【问题标题】:unable to send single message to kafka topic无法向 kafka 主题发送单个消息
【发布时间】:2017-12-15 03:53:43
【问题描述】:

我正在使用 kafka java 客户端 0.11.0 和 kafka 服务器 2.11-0.10.2.0

我的代码:

卡夫卡管理器

public class KafkaManager {

    // Single instance for producer per topic
    private static Producer<String, String> karmaProducer = null;

    /**
     * Initialize Producer
     * 
     * @throws Exception
     */
    private static void initProducer() throws Exception {
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, Constants.kafkaUrl);
        props.put(ProducerConfig.RETRIES_CONFIG, Constants.retries);
        //props.put(ProducerConfig.BATCH_SIZE_CONFIG, Constants.batchSize);
        props.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, Constants.requestTimeout);
        //props.put(ProducerConfig.LINGER_MS_CONFIG, Constants.linger);
        //props.put(ProducerConfig.ACKS_CONFIG, Constants.acks);
        //props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, Constants.bufferMemory);
        //props.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, Constants.maxBlock);
        props.put(ProducerConfig.CLIENT_ID_CONFIG, Constants.kafkaProducer);
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());

        try {
            karmaProducer = new org.apache.kafka.clients.producer.KafkaProducer<String, String>(props);
        }
        catch (Exception e) {
            throw e;
        }
    }

    /**
     * get Producer based on topic
     * 
     * @return
     * @throws Exception
     */
    public static Producer<String, String> getKarmaProducer(String topic) throws Exception {
        switch (topic) {
        case Constants.topicKarma :
            if (karmaProducer == null) {
                synchronized (KafkaProducer.class) {
                    if (karmaProducer == null) {
                        initProducer();
                    }
                }
            }
            return karmaProducer;

        default:
            return null;
        }
    }

    /**
     * Flush and close kafka producer
     * 
     * @throws Exception
     */
    public static void closeKafkaInstance() throws Exception {
        try {
            karmaProducer.flush();
            karmaProducer.close();
        } catch (Exception e) {
            throw e;
        }
    }
}

卡夫卡制作人

public class KafkaProducer {

    public void sentToKafka(String topic, String data) {
        Producer<String, String> producer = null;
        try {
            producer = KafkaManager.getKarmaProducer(topic);
            ProducerRecord<String, String> producerRecord = new ProducerRecord<String, String>(topic, data);
            producer.send(producerRecord);
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

主类

public class App {

    public static void main(String[] args) throws InterruptedException {

        System.out.println("Hello World! I am producing to stream " + Constants.topicKarma);
        String value = "google";
        KafkaProducer kafkaProducer = new KafkaProducer();
        for (int i = 1; i <= 1; i++) {
            kafkaProducer.sentToKafka(Constants.topicKarma, value + i);
            //Thread.sleep(100);
            System.out.println("Send data to producer=" + value);
            System.out.println("Send data to producer=" + value + i + " to tpoic=" +  Constants.topicKarma);
        }
    }
}

我的问题是什么:

当我的循环长度在 1000 左右(在 App 类中)时,我可以成功地将数据发送到 Kafka 主题。

但是当我的循环长度为 1 或小于 10 时,我无法将数据发送到 Kafka 主题。请注意,我没有收到任何错误。

根据我的发现,如果我想向 Kafka 主题发送一条消息,根据这个程序,我得到了成功的消息,但从未收到关于我的主题的消息。

但如果我使用 Thread.sleep(10)(正如您在我的 App 类中看到的那样,我已经对其进行了评论),那么我就成功地发送了关于我的主题的数据。

您能否解释一下为什么卡夫卡会表现出这种模棱两可的行为。

【问题讨论】:

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


    【解决方案1】:

    对 KafkaProducer.send() 的每次调用都返回一个 Future。您可以在退出之前使用最后一个 Futures 来阻塞主线程。更简单的是,您可以在发送所有消息后调用 KafkaProducer.flush(): http://kafka.apache.org/0110/javadoc/org/apache/kafka/clients/producer/KafkaProducer.html#flush()

    调用此方法可以立即发送所有缓冲的记录(即使 linger.ms 大于 0),并阻止与这些记录关联的请求完成。

    【讨论】:

      【解决方案2】:

      您正面临这个问题,因为生产者以异步方式执行发送。当您发送时,消息被放入内部缓冲区中,以便获得更大的批次,然后一次性发送更多消息。 此批处理功能配置了 batch.size 和 linger.ms,这意味着当批处理大小达到该值或经过延迟时间时发送消息。

      我在这里回复了类似的内容:Cannot produce Message when Main Thread sleep less than 1000

      即使您说“当我的循环长度约为 1000(在 App 类中)时,我也能够成功地将数据发送到 Kafka 主题。” ...但也许您看不到所有已发送的消息,因为未发送最新批次。使用较短的循环,无法及时达到上述条件,因此您在生产者有足够的时间/批量大小进行发送之前关闭应用程序。

      【讨论】:

      • 你真的想要一个只发送和关闭的应用程序吗?也许你应该让它同步工作。所以 send() 方法会返回一个 Future,您可以在其上调用 get() 方法。通过这种方式,您可以确保主线程阻塞,直到发送完成。当然吞吐量会更少。
      【解决方案3】:

      你能在退出 main 之前添加Thread.sleep(100); 吗? 如果我理解正确,那么如果您睡一小会儿,一切都会很好。如果是这种情况,则意味着您的应用程序在异步发送消息之前被终止。

      【讨论】:

      • 是的,我确认这种行为。当我使用 thread.sleep(10) 时,一切正常。但据我了解,异步方法在新线程中发送数据。这是kafka中的错误吗
      • 我不认为这是一个错误。所有线程都在退出时终止。
      • 还有一个方法重载了close(long timeout, java.util.concurrent.TimeUnit timeUnit)方法。此方法等待生产者完成所有未完成请求的发送,直至超时。尝试给出超时值。
      • 假设我有线程 t1(主类线程)和 t2(Kafka 生产者线程),如果 t1 退出,那不应该使 t2 退出
      【解决方案4】:

      将 server.properties 中的这一行从 localhost 更改为 IP 地址:

      zookeeper.connect=localhost:2181
      
      advertised.listeners=PLAINTEXT://localhost:9092
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2016-11-10
        • 1970-01-01
        • 2019-01-16
        • 2020-11-27
        • 2019-10-13
        • 2019-02-14
        相关资源
        最近更新 更多