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

如何在Kafka Streams/KTable中按多键进行GroupBy分组统计

解决Kafka Streams按多键分组并聚合到KTable的问题

要实现按syncId和status双维度分组计数,核心是创建复合分组键,将两个字段作为分组依据,再执行聚合操作。以下是具体实现步骤和代码示例:

1. 定义数据模型类

首先创建对应JSON结构的POJO,方便Kafka Streams进行序列化/反序列化:

public class HttpResponse {
    private String syncId;
    private String status;
    private String url;

    // 必须提供无参构造器,用于JSON反序列化
    public HttpResponse() {}

    // Getter和Setter方法
    public String getSyncId() { return syncId; }
    public void setSyncId(String syncId) { this.syncId = syncId; }
    public String getStatus() { return status; }
    public void setStatus(String status) { this.status = status; }
    public String getUrl() { return url; }
    public void setUrl(String url) { this.url = url; }
}

2. 构建Streams拓扑实现多键分组

通过groupBy操作将syncId和status组合为复合键,再执行计数聚合,最终得到目标KTable:

import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.Consumed;
import org.apache.kafka.streams.kstream.Grouped;
import org.apache.kafka.streams.kstream.KTable;
import org.apache.kafka.streams.kstream.Produced;
import java.util.Properties;

public class MultiKeyGroupingExample {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "multi-key-grouping-app");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());

        StreamsBuilder builder = new StreamsBuilder();

        // 从httpresponse主题读取数据,反序列化为HttpResponse对象
        var httpResponseStream = builder.stream(
            "httpresponse",
            Consumed.with(Serdes.String(), new JsonSerde<>(HttpResponse.class))
        );

        // 按syncId+status复合键分组,这里用KeyValue<String, String>承载复合键
        var groupedStream = httpResponseStream.groupBy(
            (key, value) -> new KeyValue<>(value.getSyncId(), value.getStatus()),
            Grouped.with(
                Serdes.stringSerde(),
                new JsonSerde<>(HttpResponse.class)
            )
        );

        // 执行计数聚合,得到KTable:键为(syncId, status),值为对应计数
        KTable<KeyValue<String, String>, Long> countTable = groupedStream.count();

        // 可选:将结果转换为结构化格式输出到新主题(方便查看或下游消费)
        countTable.toStream().map(
            (compositeKey, count) -> new KeyValue<>(
                compositeKey.key,
                String.format("%s,%d", compositeKey.value, count)
            )
        ).to("httpresponse-counts", Produced.with(Serdes.String(), Serdes.String()));

        KafkaStreams streams = new KafkaStreams(builder.build(), props);
        streams.start();

        // JVM关闭钩子
        Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
    }
}

关键说明:

  • 复合键设计:示例用KeyValue<String, String>组合syncId和status,你也可以自定义包含这两个字段的POJO作为键,只需为其实现对应的Serde(序列化/反序列化器)。
  • 分组逻辑差异:groupBy允许你从原数据中提取新的分组键(这里就是双字段组合),而groupByKey只能基于原消息的键分组——这是你之前仅按syncId分组结果不符合预期的核心原因。
  • KTable持久化:聚合后的countTable本身就是持久化的KTable,Kafka Streams会将其状态存储在默认的RocksDB中;如果需要持久化到外部存储或其他主题,可通过to()方法输出。

3. 自定义JsonSerde实现

如果没有现成的JSON Serde,可基于Jackson实现通用序列化器:

import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.kafka.common.serialization.Deserializer;
import org.apache.kafka.common.serialization.Serde;
import org.apache.kafka.common.serialization.Serializer;
import java.util.Map;

public class JsonSerde<T> implements Serde<T> {
    private final ObjectMapper objectMapper = new ObjectMapper();
    private final Class<T> type;

    public JsonSerde(Class<T> type) {
        this.type = type;
    }

    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {}

    @Override
    public void close() {}

    @Override
    public Serializer<T> serializer() {
        return (topic, data) -> {
            try {
                return objectMapper.writeValueAsBytes(data);
            } catch (Exception e) {
                throw new RuntimeException("JSON序列化失败", e);
            }
        };
    }

    @Override
    public Deserializer<T> deserializer() {
        return (topic, data) -> {
            try {
                return objectMapper.readValue(data, type);
            } catch (Exception e) {
                throw new RuntimeException("JSON反序列化失败", e);
            }
        };
    }
}

内容的提问来源于stack exchange,提问作者voila

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 03:20:14