【问题标题】:How to know howmany data consumed from kafka queue and what are data existing inside kafka queue topic?如何知道从 kafka 队列中消耗了多少数据以及 kafka 队列主题中存在哪些数据?
【发布时间】:2016-10-26 20:44:34
【问题描述】:

我在“topicDemo”Kafka 中生成和使用数据,以启用分布式数据并并行处理这些数据。

但在实时场景中需要监控当前队列中有多少数据(topicDemo)以及从(topicDemo)消耗了多少数据。

是否有任何 kafka API 可用于提供这些详细信息?

这是我正在生成数据的代码

  // create instance for properties to access producer configs
        Properties props = new Properties();

        // props.put("serializer.class",
        // "kafka.serializer.StringEncoder");
        props.put("bootstrap.servers", "localhost:9092");
        props.put("metadata.broker.list", "localhost:9092");

        props.put("producer.type", "async");
        props.put("batch.size", "500");
        props.put("compression.codec", "1");
        props.put("compression.topic", "topicdemo");
        // props.put("key.serializer",
        // "org.apache.kafka.common.serialization.StringSerializer");
        props.put("key.serializer", "org.apache.kafka.common.serialization.IntegerSerializer");
        props.put("value.serializer", "org.apache.kafka.common.serialization.ByteArraySerializer");

        org.apache.kafka.clients.producer.Producer<Integer, byte[]> producer = new KafkaProducer<Integer, byte[]>(
                props);

            producer.send(new ProducerRecord<Integer, byte[]>("topicdemo", resume.getBytes()));
        producer.close();

这是我正在使用数据的代码

  String topicName = "topicDemo";

    Properties props = new Properties();

    props.put("bootstrap.servers", "localhost:9092");
    props.put("group.id", "test");
    props.put("enable.auto.commit", "true");
    props.put("auto.commit.interval.ms", "1000");
    props.put("session.timeout.ms", "30000");
    props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
    props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");

    KafkaConsumer<String, String> consumer = new KafkaConsumer<String, String>(props);

    // Kafka Consumer subscribes list of topics here.
    consumer.subscribe(Arrays.asList(topicName));

    try {
        while (true) {
            ConsumerRecords<String, String> records = consumer.poll(5);
            for (ConsumerRecord<String, String> record : records) {
                Consumer cons = new Consumer();
                if (cons.SaveDetails(record.value()) == 1) {
                    consumer.commitSync();
                }
            }
        }
    } catch (Exception e) {
        e.printStackTrace();
    } finally {
        consumer.close();
    }

【问题讨论】:

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


    【解决方案1】:

    执行下面的命令来检查你为每个分区产生了多少消息:

    bin/kafka-run-class.sh kafka.tools.GetOffsetShell --broker-list host:port --topic topic_name --time -1

    执行以下命令查看消费者已消费和落后的消息数量:

    bin/kafka-consumer-groups.sh --bootstrap-server broker1:9092 --describe --group test-consumer-group(老消费者)

    bin/kafka-consumer-groups.sh --bootstrap-server broker1:9092 --describe --group test-consumer-group --new-consumer(新消费者)

    【讨论】:

    • 有没有kafka提供的API让我可以用java实现和准备监控工具?
    • 您可以调用 GetOffsetShell.main(args) 来模拟第一个命令;对于第二个命令,尝试调用 ConsumerGroupCommand.main(args)
    猜你喜欢
    • 2016-03-23
    • 1970-01-01
    • 2014-09-20
    • 1970-01-01
    • 1970-01-01
    • 2018-11-09
    • 2018-03-17
    • 2019-08-18
    • 2019-02-20
    相关资源
    最近更新 更多