【问题标题】:Get all kafka messages in a queue and stop streaming in java获取队列中的所有kafka消息并在java中停止流式传输
【发布时间】:2016-08-29 00:12:28
【问题描述】:

我需要在晚上执行一个作业,它将获取 kafka 队列中的所有消息并使用它们执行一个进程。我能够收到消息,但 kafka 流正在等待更多消息,我无法继续我的流程。我有以下代码:

...
private ConsumerConnector consumerConnector;
private final static String TOPIC = "test";

public MessageStreamConsumer() {
        Properties properties = new Properties();
        properties.put("zookeeper.connect", "localhost:2181");
        properties.put("group.id", "test-group");
        ConsumerConfig consumerConfig = new ConsumerConfig(properties);
        consumerConnector = Consumer.createJavaConsumerConnector(consumerConfig);
    }
public List<String> getMessages() {
                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);
                ConsumerIterator<byte[], byte[]> it = stream.iterator();
                List<String> messages = new ArrayList<>();
                while (it.hasNext())
                    messages.add(new String(it.next().message()));
                return messages;
            }

代码能够获取消息,但是当它处理最后一条消息时,它会留在行中:

 while (it.hasNext())

问题是,我怎样才能从 kafka 获取所有消息,停止流并继续我的其他任务。

希望你能帮到我

谢谢

【问题讨论】:

  • 但我不认为这样做的最佳做法是等到抛出异常。如果我的进程花费的时间超过了配置的超时时间怎么办
  • 你不应该在这里使用直接的 Kafka Consumer,而不是 KafkaStream 吗?流自然会保持活力。

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


【解决方案1】:

好像kafka流不支持从头消费。
您可以创建一个原生 kafka 消费者并将auto.offset.reset 设置为最早,然后它将从头开始消费消息。

【讨论】:

    【解决方案2】:

    这样的事情可能会奏效。基本上,这个想法是使用 Kafka Consumer 并轮询,直到你得到一些记录,然后在你得到一个空批次时停止。

    package kafka.examples;
    
    import java.text.DateFormat;
    import java.text.SimpleDateFormat;
    import java.util.Calendar;
    import java.util.Collections;
    import java.util.Date;
    import java.util.Properties;
    import java.util.concurrent.CountDownLatch;
    import java.util.concurrent.atomic.AtomicBoolean;
    
    import org.apache.kafka.clients.consumer.ConsumerRecord;
    import org.apache.kafka.clients.consumer.ConsumerRecords;
    import org.apache.kafka.clients.consumer.KafkaConsumer;
    
    
    public class Consumer1 extends Thread
    {
        private final KafkaConsumer<Integer, String> consumer;
        private final String topic;
        private final DateFormat df;
        private final String logTag;
        private boolean noMoreData = false;
        private boolean gotData = false;
        private int messagesReceived = 0;
        AtomicBoolean isRunning = new AtomicBoolean(true);
        CountDownLatch shutdownLatch = new CountDownLatch(1);
    
        public Consumer1(Properties props)
        {
            logTag = "Consumer1";
    
            consumer = new KafkaConsumer<>(props);
            this.topic = props.getProperty("topic");
            this.df = new SimpleDateFormat("HH:mm:ss");
    
            consumer.subscribe(Collections.singletonList(this.topic));
        }
    
        public void getMessages() {
            System.out.println("Getting messages...");
            while (noMoreData == false) {
                //System.out.println(logTag + ": Doing work...");
    
                ConsumerRecords<Integer, String> records = consumer.poll(1000);
                Date now = Calendar.getInstance().getTime();
                int recordsCount = records.count();
                messagesReceived += recordsCount;
                System.out.println("recordsCount: " + recordsCount);
                if (recordsCount > 0) {
                   gotData = true;
                }
    
                if (gotData && recordsCount == 0) {
                    noMoreData = true;
                }
    
                for (ConsumerRecord<Integer, String> record : records) {
                    int kafkaKey = record.key();
                    String kafkaValue = record.value();
                    System.out.println(this.df.format(now) + " " + logTag + ":" +
                            " Received: {" + kafkaKey + ":" + kafkaValue + "}" +
                            ", partition(" + record.partition() + ")" +
                            ", offset(" + record.offset() + ")");
                }
            }
            System.out.println("Received " + messagesReceived + " messages");
        }
    
        public void processMessages() {
            System.out.println("Processing messages...");
        }
    
        public void run() {
            getMessages();
            processMessages();
        }
    }
    

    【讨论】:

      【解决方案3】:

      我目前正在使用 Kafka 0.10.0.1 进行开发,发现关于使用消费者属性 auto.offset.reset 的混合信息,所以我做了一些实验来弄清楚实际发生了什么。

      基于这些,我现在是这样理解的:当你设置属性时:

      auto.offset.reset=earliest
      

      这将消费者定位到分配的分区中的第一个可用消息(当分区上没有提交时)或者它将消费者定位在最后提交的分区偏移量(请注意,您始终提交最后一次读取偏移量 + 1否则您将在每次重启消费者时重新阅读最后提交的消息)

      或者您不设置 auto.offset.reset,这意味着将使用“最新”的默认值。

      在这种情况下,您在连接消费者时不会收到任何旧消息 - 只会收到连接消费者后发布到主题的消息。

      作为结论 - 如果您想确保接收某个主题和分配的分区的所有可用消息,您必须调用 seekToBeginning()。

      似乎建议首先调用 poll(0L) 以确保您的消费者获得分配的分区(或在 ConsumerRebalanceListener 中实现您的代码!),然后将每个分配的分区寻找到“开始”:

      kafkaConsumer.poll(0L);
      kafkaConsumer.seekToBeginning(kafkaConsumer.assignment());
      

      【讨论】:

        猜你喜欢
        • 2020-08-15
        • 2015-01-26
        • 2013-12-13
        • 2018-08-10
        • 1970-01-01
        • 1970-01-01
        • 2017-09-15
        • 1970-01-01
        • 2013-03-24
        相关资源
        最近更新 更多