【问题标题】:Kafka--Consumer reading exactly half the stream卡夫卡——消费者准确地阅读了流的一半
【发布时间】:2016-01-22 23:29:02
【问题描述】:

我正在使用以下代码来读取我的主题数据,即“sha-test2”,但它正在读取完全替代的代码行,即 20 行中的 10 行。 但是当我运行控制台时,它显示了所有 20 行。 IE 。 bin/kafka-console-consumer.sh --zookeeper localhost:2181 --topic sha-test2 --from-beginning

我哪里错了?非常感谢您的帮助。

public class KafkaTestConsumer extends  Thread {
    //final static String clientId = "SimpleConsumerDemoClient";
    final static String TOPIC = "sha-test2";
    ConsumerConnector consumerConnector;

    public static void main(String[] argv) throws   
     UnsupportedEncodingException {
        KafkaTestConsumer helloKafkaConsumer = new KafkaTestConsumer();
        helloKafkaConsumer.start();
    }
    public KafkaTestConsumer(){
        Properties properties = new Properties();
        properties.put("zookeeper.connect","172.23.32.35:2181");
        properties.put("group.id","test-group");
        ConsumerConfig consumerConfig = new ConsumerConfig(properties);
        consumerConnector = 
         Consumer.createJavaConsumerConnector(consumerConfig);
    }


    @Override
    public void run() {
        Map<String, Integer> topicCountMap = new HashMap<String, Integer>();
        topicCountMap.put(TOPIC, new Integer(1));
        Map<String, List<KafkaStream<byte[], byte[]>>> consumerMap =  
         consumerConnector.createMessageStreams(topicCountMap);
        KafkaStream<byte[], byte[]> stream =  consumerMap.get(TOPIC).get(0);
        System.out.println("consumerMap : \n " + consumerMap.toString() );
        ConsumerIterator<byte[], byte[]> it = stream.iterator();

       System.out.println("run started");
        while(it.hasNext()){
            System.out.println(new String(it.next().message()));
        }
}

Thank you.
~Shyam

【问题讨论】:

    标签: regex apache-kafka hadoop-streaming kafka-consumer-api


    【解决方案1】:

    问题出在这一行:

    topicCountMap.put(TOPIC, new Integer(1));
    

    您告诉consumerConnector 为您的主题创建一个消费者线程,但该主题(显然)有两个分区。 "test-group" 组中的消费者线程数应等于或大于分区数,否则该组将无法读取某些分区,这正是您的情况。

    请查看this example,其中线程数是通过命令行参数设置的。

    或者,您可以从 /brokers/topics/your_topic_name/partitions 节点下的 Zookeeper 中读取存储元数据的分区的确切数量。

    【讨论】:

      【解决方案2】:

      您的代码看起来非常好。这看起来像一个偏移问题。高级消费者将其偏移量存储在 Zookeeper 中。

      在您的情况下,这可能会发生:- 1.你在kafka里放了10条消息 2. 你运行了消费者代码,它成功读取了所有 10 条消息。此外,消费者在 zookeeper 中将消耗的偏移量更新为 10。 3. 你阻止你的消费者。 4.你又给kafka发了10条消息 5. 你再次启动消费者代码。它只读取最后 10 条消息,而不是之前推送的 10 条消息,因为当您重新启动消费者时,它会检查 zookeeper 以找出从哪个偏移量恢复消费。

      尝试使用不同的组 id 重新运行您的消费者,或者尝试从 zookeeper 中删除偏移量。它应该可以正常工作。

       properties.put("group.id","test-group420");
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2021-04-02
        • 1970-01-01
        • 2019-07-03
        • 2018-05-05
        • 2021-08-22
        • 1970-01-01
        相关资源
        最近更新 更多