You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何为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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.07 06:24:51