【问题标题】:Writing a kafka topic when consumer is down当消费者情绪低落时写一个 kafka 主题
【发布时间】:2016-01-10 05:32:47
【问题描述】:

我试图整合 kafka-storm。我只是从几个例子开始。

我能够从 GitHub 运行示例。接下来我尝试在 Eclipse 中编写一个 Producer 类,以使用 KAFKA PRODUCER API 将消息发布到 kafka 主题。

场景1:

当我的消费者外壳使用说主题测试运行时,我运行我的生产者类。我可以看到我的消费者外壳以及所有已发布的消息。

场景2

我还没有启动我的消费者外壳(假设消费者已关闭)。我经营我的制作人课程。消息正在发布到 kafka。

现在如果消息已发布,现在如果我启动消费者 shell,则在停机后,它不会读取已发布的消息主题。

为什么?我想它会维护主题消费的日志。不应该是在看消息吗?

有什么配置参数需要提一下吗?

Properties props = new Properties();
props.put("metadata.broker.list", "localhost:9092");
props.put("zk.connect", "localhost:2181");
props.put("serializer.class", "kafka.serializer.StringEncoder");
props.put("request.required.acks", "1");

ProducerConfig config = new ProducerConfig(props);
Producer<String, String> producer = new Producer<String, String>(config);

        for ( int nEvents=0; nEvents<events;nEvents++)
        {
          String ip="192.168.2."+rnd.nextInt(255);
          String msg=getNextTradeData(); // Class to generate data
          KeyedMessage<String,String> data=new KeyedMessage<String, String>("TradeFrequency",ip,msg);
          Thread.sleep(100);
          System.out.println(msg);
          producer.send(data);  

        }
producer.close();

}

或者我需要做些什么来改变消费者。我正在使用包中提供的 consumer-shell,并使用

启动它
  bin/kafka-console-consumer.sh --zookeeper localhost:2181 --topic first-topic

【问题讨论】:

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


    【解决方案1】:

    当您启动kafka-console-consume 时,它将从当前偏移量读取。这是NOW() 的偏移量,而不是过去的偏移量。

    要查看消息是否已发布,您有两种选择:

    1. 使用--from-beginning 选项,从主题的开头开始阅读

      bin/kafka-console-consumer.sh --zookeeper localhost:2181 --topic first-topic --from-beginning

    2. 使用--consumer.config选项在zookeeper/kafka中保持console-consumer的状态

      bin/kafka-console-consumer.sh --zookeeper localhost:2181 --topic first-topic --consumer.config /home/sql-injection/consumer-config.txt

    根据这个nice page,您需要在消费者配置上的参数是:consumer.id,client.id。

    【讨论】:

    • 谢谢,在这种情况下它对我有用。现在在这种情况下,每次它都会从我想的日志的 bignning 开始(尽管我仍然需要检查)。但是,如果我想从它离开的任何地方读取,或者可以说(最后一个偏移量读取+1)。我的方法是什么?
    • @NishantM 使用消费者配置选项,这应该允许您从最后一个偏移量恢复
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-05-31
    • 2017-01-26
    • 1970-01-01
    • 1970-01-01
    • 2020-10-12
    • 2018-12-31
    • 2017-10-17
    相关资源
    最近更新 更多