kafka消费者代码不运行完全

kafka consumer code is not running completely

我正在尝试一个简单的 kafka 生产者消费者客户端 api,我的生产者 class 工作正常,因为我可以从控制台看到消费者中的消息但是当我是 运行 消费者时代码什么都没有显示,我不知道是什么问题或我哪里做错了

这是生产者代码是 -

package com.app.test;
import java.util.Properties;
import kafka.javaapi.producer.Producer;
import kafka.producer.KeyedMessage;
import kafka.producer.ProducerConfig;

public class SimpleProducer {
    private static Producer<Integer, String> producer;
    private final Properties props = new Properties();public    SimpleProducer()
    {
      props.put("metadata.broker.list", "localhost:9092");
      props.put("serializer.class", "kafka.serializer.StringEncoder");
      props.put("request.required.acks", "1");
      producer = new Producer<Integer, String>(new ProducerConfig(props));
    }
public static void main(String[] args) {
    SimpleProducer sp = new SimpleProducer();
    String topic = "test4";
    String messageStr = "hello";
    KeyedMessage<Integer, String> data = new KeyedMessage<Integer, String>(topic, messageStr);
    System.out.println("producer : "+producer);
    producer.send(data);
    producer.close();
  }
}

消费者 class 是 -

package com.app.test;

import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Properties;

import kafka.consumer.Consumer;
import kafka.consumer.ConsumerConfig;
import kafka.consumer.ConsumerIterator;
import kafka.consumer.KafkaStream;
import kafka.javaapi.consumer.ConsumerConnector;

public class SimpleHLConsumer {

  private final ConsumerConnector consumer;
  private final String topic;

  public SimpleHLConsumer(String zookeeper, String groupId, String topic) {
      Properties props = new Properties();
    props.put("zookeeper.connect", zookeeper);
    props.put("group.id", groupId);
    props.put("zookeeper.session.timeout.ms", "500");
    props.put("zookeeper.sync.time.ms", "250");
    props.put("auto.commit.interval.ms", "1000");

    consumer = Consumer.createJavaConsumerConnector(
    new ConsumerConfig(props));
    this.topic = topic;
  }



  public void testConsumer() {


    Map<String, Integer> topicCount = new HashMap<String, Integer>();
        // Define single thread for topic
    topicCount.put(topic, new Integer(1));      

    System.out.println("check1");

    Map<String, List<KafkaStream<byte[], byte[]>>> consumerStreams = consumer.createMessageStreams(topicCount);

    System.out.println("check2");
    List<KafkaStream<byte[], byte[]>> streams = consumerStreams.get(topic);

   for (KafkaStream stream : streams) {
        System.out.println("test----");
        System.out.println("test----"+stream.toString());
        ConsumerIterator<byte[], byte[]> consumerIte = stream.iterator();
        while (consumerIte.hasNext()) {
            try {
                System.out.println("Message from Single Topic: " + new String(consumerIte.next().message(), "UTF-8"));
            } catch (UnsupportedEncodingException e) {
                // TODO Auto-generated catch block
                e.printStackTrace();
            }
        }
    }
    if (consumer != null)
      consumer.shutdown();
  }



  public static void main(String[] args) {

    String topic = "test4";
SimpleHLConsumer simpleHLConsumer = new SimpleHLConsumer("localhost:2181", "testgroup", topic);
    simpleHLConsumer.testConsumer();
  }
}

为了检查,我在 testConsumer() 方法中使用 sysout 进行了 2 次检查,因此虽然 运行 仅显示 check1,即代码未到达 check2,我认为 [=13 存在一些问题=],那是什么原因,我该如何解决呢?

代码没有问题。

只需构建包含 kafka-version(您正在使用的 kafka 版本)的 jar 和 运行 它并尝试从生产者控制台发送消息。

希望对您有所帮助。