Kafka Java API获取非compacted topic总消息数

目前Kafka并没有提供直接的工具来帮助我们获取某个topic的当前总消息数,需要我们自行写程序来实现。下列代码可以实现这一功能,特此记录一下:

1/** 2 * 获取某个topic的当前消息数 3 * Java 8+ only 4 * 5 * @param topic 6 * @param brokerList 7 * @return 8 */ 9 public static long totalMessageCount(String topic, String brokerList) { 10 Properties props = new Properties(); 11 props.put("bootstrap.servers", brokerList); 12 props.put("group.id", "test-group"); 13 props.put("enable.auto.commit", "false"); 14 props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); 15 props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); 16 17 try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) { 18 List<TopicPartition> tps = Optional.ofNullable(consumer.partitionsFor(topic)) 19 .orElse(Collections.emptyList()) 20 .stream() 21 .map(info -> new TopicPartition(info.topic(), info.partition())) 22 .collect(Collectors.toList()); 23 Map<TopicPartition, Long> beginOffsets = consumer.beginningOffsets(tps); 24 Map<TopicPartition, Long> endOffsets = consumer.endOffsets(tps); 25 26 return tps.stream().mapToLong(tp -> endOffsets.get(tp) - beginOffsets.get(tp)).sum(); 27 } 28 }
点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

Java日期时间API系列31

  时间戳是指格林威治时间1970年01月01日00时00分00秒起至现在的总毫秒数,是所有时间的基础,其他时间可以通过时间戳转换得到。Java中本来已经有相关获取时间戳的方法,Java8后增加新的类Instant等专用于处理时间戳问题。 1获取时间戳的方法和性能对比1.1获取时间戳方法Java8以前