如何为Kafka Streams各流单独统计记录数及获取偏移量?
Kafka Streams 分输入主题统计处理指标
当前拓扑与需求背景
我正在使用org.apache.kafka.streams.KafkaStreams,拓扑结构如下:
StreamsBuilder builder = new StreamsBuilder(); builder.stream("input-topic1") .mapValues((readOnlyKey, value) -> value.toUpperCase()) .to("output-topic1"); builder.stream("input-topic2") .mapValues((readOnlyKey, value) -> value.toUpperCase()) .to("output-topic2");
默认情况下Kafka Streams每两分钟会输出类似日志:
Processed 14 total records, ran 0 punctuators, and committed 11 total tasks since the last update
但我需要更清晰地掌握每个输入主题的消息流入情况,希望获取各流的详细指标(而非仅总处理记录数),具体需求包括:
- 为每个流命名,分别提取
builder.stream("input-topic1")和builder.stream("input-topic2")的消费偏移量 - 统计某时间段内各流处理的记录数
已尝试方案及问题
- 使用
.peek配合静态变量统计:属于不良实践,存在线程安全风险,不符合Kafka Streams状态管理规范 - 查看KafkaStreams metrics:未找到直接对应分主题的处理记录/偏移量指标
- 每条消息记录日志:会产生大量日志,不可行
可行解决方案
1. 利用Kafka Streams内置Metrics按主题拆分统计
Kafka Streams的Metrics包含分主题、分分区的处理指标,需通过维度筛选获取:
- 记录处理数:指标名称为
stream-processor-node-metrics,通过topic维度过滤可获取对应输入主题的累计处理记录数 - 偏移量进度:指标名称为
stream-fetch-metrics,通过topic、partition维度可获取各输入主题分区的当前消费偏移量、滞后量
示例代码:
KafkaStreams streams = new KafkaStreams(topology, config); streams.start(); // 每2分钟获取一次分主题统计数据 ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(); executor.scheduleAtFixedRate(() -> { MetricRegistry metricRegistry = streams.metrics(); metricRegistry.getMetrics().forEach((metricName, metric) -> { // 筛选input-topic1的累计处理记录数 if ("stream-processor-node-metrics".equals(metricName.getName()) && "input-topic1".equals(metricName.getTags().get("topic")) && "records-processed-total".equals(metricName.getTags().get("metric-name"))) { System.out.printf("input-topic1 累计处理记录数: %d%n", (long) metric.metricValue()); } // 同理可筛选input-topic2的指标 }); }, 0, 2, TimeUnit.MINUTES);
2. 自定义Processor结合状态存储统计
通过自定义Processor,搭配Kafka Streams的持久化状态存储,安全统计各主题的周期处理记录数,避免静态变量的线程安全问题:
// 自定义统计Processor class TopicCountProcessor implements Processor<String, String> { private ProcessorContext context; private KeyValueStore<String, Long> countStore; @Override public void init(ProcessorContext context) { this.context = context; // 获取预先定义的状态存储 countStore = (KeyValueStore<String, Long>) context.getStateStore("topic-count-store"); // 每2分钟输出统计结果并重置计数 context.schedule(Duration.ofMinutes(2), PunctuationType.WALL_CLOCK_TIME, timestamp -> { countStore.all().forEachRemaining(entry -> { System.out.printf("主题 %s 最近2分钟处理记录数: %d%n", entry.key, entry.value); countStore.put(entry.key, 0L); }); }); } @Override public void process(String key, String value) { // 通过上下文获取当前消息所属主题 String topic = context.topic(); Long currentCount = countStore.get(topic); countStore.put(topic, currentCount == null ? 1 : currentCount + 1); // 向下游传递消息 context.forward(key, value); } @Override public void close() {} } // 集成到拓扑中 StreamsBuilder builder = new StreamsBuilder(); // 创建持久化状态存储 builder.addStateStore(Stores.keyValueStoreBuilder( Stores.persistentKeyValueStore("topic-count-store"), Serdes.String(), Serdes.Long() )); // 为每个输入流添加统计Processor builder.stream("input-topic1") .process(() -> new TopicCountProcessor(), "topic-count-store") .mapValues((readOnlyKey, value) -> value.toUpperCase()) .to("output-topic1"); builder.stream("input-topic2") .process(() -> new TopicCountProcessor(), "topic-count-store") .mapValues((readOnlyKey, value) -> value.toUpperCase()) .to("output-topic2");
3. 通过AdminClient查询消费偏移量
如果需要直接获取各输入主题的消费偏移量,可使用Kafka AdminClient查询Kafka Streams对应的消费者组偏移量(消费者组名由application.id配置):
AdminClient adminClient = AdminClient.create(config); String appId = config.getProperty(StreamsConfig.APPLICATION_ID_CONFIG); // 查询消费者组的分区分配与偏移量 ConsumerGroupDescription groupDesc = adminClient.describeConsumerGroups(Collections.singletonList(appId)) .all() .get() .values() .iterator() .next(); groupDesc.members().forEach(member -> { member.assignment().topicPartitions().forEach(tp -> { if (tp.topic().matches("input-topic\\d+")) { adminClient.listConsumerGroupOffsets(appId) .partitionsToOffsetAndMetadata() .get() .forEach((partition, offsetMeta) -> { if (partition.equals(tp)) { System.out.printf("主题 %s 分区 %d 当前消费偏移量: %d%n", tp.topic(), tp.partition(), offsetMeta.offset()); } }); } }); });
内容的提问来源于stack exchange,提问作者Vytautas Šerėnas
相关产品推荐
相关产品推荐

