【问题标题】:How to get Kafka Producer messages count如何获取 Kafka Producer 消息计数
【发布时间】:2020-06-19 05:15:56
【问题描述】:

我使用以下代码创建了一个生产者,它产生了大约 2000 条消息。

public class ProducerDemoWithCallback {

    public static void main(String[] args) {

        final Logger logger = LoggerFactory.getLogger(ProducerDemoWithCallback.class);

String bootstrapServers = "localhost:9092";
    Properties properties = new Properties();
    properties.setProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
    properties.setProperty(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
    properties.setProperty(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());

    // create the producer
    KafkaProducer<String, String> producer = new KafkaProducer<String, String>(properties);


    for (int i=0; i<2000; i++ ) {
        // create a producer record
        ProducerRecord<String, String> record =
                new ProducerRecord<String, String>("TwitterProducer", "Hello World " + Integer.toString(i));

        // send data - asynchronous
        producer.send(record, new Callback() {
            public void onCompletion(RecordMetadata recordMetadata, Exception e) {
                // executes every time a record is successfully sent or an exception is thrown
                if (e == null) {
                    // the record was successfully sent
                    logger .info("Received new metadata. \n" +
                            "Topic:" + recordMetadata.topic() + "\n" +
                            "Partition: " + recordMetadata.partition() + "\n" +
                            "Offset: " + recordMetadata.offset() + "\n" +
                            "Timestamp: " + recordMetadata.timestamp());


                } else {
                    logger .error("Error while producing", e);
                }
            }
        });
    }

    // flush data
    producer.flush();
    // flush and close producer
    producer.close();
  }
}

我想计算这些消息并获得 int 值。 我使用此命令并且它有效,但我正在尝试使用代码获取此计数。

"bin/kafka-run-class.sh kafka.tools.GetOffsetShell --broker-list localhost:9092 --topic TwitterProducer --time -1"

结果是

- TwitterProducer:0:2000

我以编程方式执行相同操作的代码如下所示,但我不确定这是否是获取计数的正确方法:

 int valueCount = (int) recordMetadata.offset();
 System.out.println("Offset value " + valueCount); 

有人可以帮助我使用代码获取 Kafka 消息偏移值的计数。

【问题讨论】:

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


    【解决方案1】:

    你可以看看GetOffsetShell的实现细节。

    这是一个用 Java 重写的简化代码:

    import org.apache.kafka.clients.consumer.ConsumerConfig;
    import org.apache.kafka.clients.consumer.KafkaConsumer;
    import org.apache.kafka.common.TopicPartition;
    import org.apache.kafka.common.serialization.StringDeserializer;
    
    import java.util.*;
    import java.util.stream.Collectors;
    
    public class GetOffsetCommand {
    
        private static final Set<String> TopicNames = new HashSet<>();
    
        static {
            TopicNames.add("my-topic");
            TopicNames.add("not-my-topic");
        }
    
        public static void main(String[] args) {
            TopicNames.forEach(topicName -> {
                final Map<TopicPartition, Long> offsets = getOffsets(topicName);
    
                new ArrayList<>(offsets.entrySet()).forEach(System.out::println);
                System.out.println(topicName + ":" + offsets.values().stream().reduce(0L, Long::sum));
            });
        }
    
        private static Map<TopicPartition, Long> getOffsets(String topicName) {
            final KafkaConsumer<String, String> consumer = makeKafkaConsumer();
            final List<TopicPartition> partitions = listTopicPartitions(consumer, topicName);
            return consumer.endOffsets(partitions);
        }
    
        private static KafkaConsumer<String, String> makeKafkaConsumer() {
            final Properties props = new Properties();
            props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
            props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
            props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
            props.put(ConsumerConfig.GROUP_ID_CONFIG, "get-offset-command");
    
            return new KafkaConsumer<>(props);
        }
    
        private static List<TopicPartition> listTopicPartitions(KafkaConsumer<String, String> consumer, String topicName) {
            return consumer.listTopics().entrySet().stream()
                    .filter(t -> topicName.equals(t.getKey()))
                    .flatMap(t -> t.getValue().stream())
                    .map(p -> new TopicPartition(p.topic(), p.partition()))
                    .collect(Collectors.toList());
        }
    }
    

    它为每个主题的分区和总和(消息总数)生成偏移量,例如:

    my-topic-0=184
    my-topic-2=187
    my-topic-4=189
    my-topic-1=196
    my-topic-3=243
    my-topic:999
    

    【讨论】:

    • 效果很好。它的实现是使用 consumerConfig 完成的。 ProdcerConfig 不可能吗?
    • 否,因为它使用的是 Consumer API。无法从 Producer API 检索偏移元数据(因为 Kafka 中的生产者基本上不关心偏移量)。
    • 最后的澄清。如果有多个主题,例如假设我有 3 个主题,我想计算每个主题。
    • 您可以针对不同的主题名称多次运行代码或重构 .filter(t -&gt; TopicName.equals(t.getKey())) 以与多个主题名称进行比较,例如.filter(t -&gt; TopicNamesSet.contains(t.getKey())) 其中TopicNamesSet 是一个Set,包含所有所需的主题名称。
    • 谢谢。它有效地工作。我想在上面的代码中添加最后一个功能来完成它。我应该将此作为一个单独的问题提出吗?代码返回这样的结果 my-topic-0 =2000 not-my-topic-0 =2000 我试图比较 set 的这两个元素的值,但没有得到正确的结果。我在 SO 和互联网上看到的只是两组的比较。但我试图比较这两个元素在集合中的值。像这样 if (my-topic.count == not-my-topic.count) { System.out.println("它们是相等的"); }
    【解决方案2】:

    您为什么要获得该值?如果你分享更多关于目的的细节,我可以给你更多的好建议。

    对于您的最后一个问题,使用偏移值获取消息计数不是正确的方法。如果你的主题有一个分区,生产者是一个,你可以使用它。您需要考虑该主题有几个分区。

    如果要获取每个生产者的消息数量,可以在回调函数onCompletion()中统计

    或者您可以像这样使用消费者客户端获取最后一个偏移量:

    Properties props = new Properties();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "your-brokers");
    props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
    props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
    
    Consumer<Long, String> consumer = new KafkaConsumer<>(props);
    consumer.subscribe(Collections.singletonList("topic_name");
    
    Collection<TopicPartition> partitions = consumer.assignment();
    
    consumer.seekToEnd(partitions);
    
    for(TopicPartition tp: partitions) {
        long offsetPosition = consumer.position(tp);
    }
    

    【讨论】:

    • 生产者将生成消息,这些消息将存储在应用程序平均数据库中。然后我需要比较结果。从数据库方面,我可以数数。但在制作人方面,我不知道如何处理它。有没有可能我将此计数作为整数,以便我可以与数据库计数值进行比较?你能给出一些想法或例子吗?
    • 好的,我认为您可以使用我的第二个答案获得最后生成的偏移位置。然后您可以将其与您的数据库进行比较。也许这对您也有帮助:stackoverflow.com/questions/38428196/…
    猜你喜欢
    • 2018-01-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-04-22
    • 2021-11-09
    • 1970-01-01
    相关资源
    最近更新 更多