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

如何用Kafka Streams和窗口功能实现时间窗口内的计数?

嘿,很高兴能帮到你!既然你已经能成功读取Kafka主题并打印消息了,那接下来实现按ID分组+时间窗口内计数的需求其实很清晰,咱们一步步来拆解实现:

步骤1:定义消息实体类(POJO)

首先,为了更方便地处理结构化的消息,咱们先定义一个和消息格式对应的Java实体类,这样Kafka Streams可以直接完成消息的序列化/反序列化:

import com.fasterxml.jackson.annotation.JsonProperty;

public class Message {
    @JsonProperty("ID")
    private String id;
    
    @JsonProperty("time")
    private String time;

    // 必须保留无参构造函数,用于JSON反序列化
    public Message() {}

    // 带参构造、getter和setter方法
    public Message(String id, String time) {
        this.id = id;
        this.time = time;
    }

    public String getId() {
        return id;
    }

    public void setId(String id) {
        this.id = id;
    }

    public String getTime() {
        return time;
    }

    public void setTime(String time) {
        this.time = time;
    }
}
步骤2:配置JSON序列化/反序列化器(Serde)

因为你的消息是JSON格式,咱们需要为Message类配置对应的Serde,让Kafka Streams能正确解析和生成消息。可以直接用Kafka Streams提供的JSON Serde(需要引入kafka-streams-json依赖):

import org.apache.kafka.common.serialization.Serde;
import org.apache.kafka.common.serialization.Serdes;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.kafka.streams.kstream.Serdes;

// 构建Message的JSON Serde
Serde<Message> messageSerde = Serdes.serdeFrom(
    new JsonSerializer<>(new ObjectMapper()),
    new JsonDeserializer<>(Message.class, new ObjectMapper())
);
步骤3:构建流处理拓扑(核心逻辑)

接下来就是核心的分组和窗口计数逻辑了。这里需要注意:如果你的消息原始Key不是ID,就需要用groupBy指定以消息中的ID字段作为分组键,然后定义时间窗口并完成计数:

import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.kstream.*;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsConfig;
import java.time.Duration;
import java.util.Properties;

public class KafkaStreamsIdCountExample {
    public static void main(String[] args) {
        // 初始化Streams配置(替换成你的实际集群地址等配置)
        StreamsConfig streamsConfig = getStreamsConfig();

        StreamsBuilder builder = new StreamsBuilder();

        // 从输入主题读取消息
        KStream<String, Message> inputStream = builder.stream(
            "your-input-topic",
            Consumed.with(Serdes.String(), messageSerde)
        );

        // 按消息中的ID字段分组
        KGroupedStream<String, Message> groupedById = inputStream.groupBy(
            (originalKey, message) -> message.getId(),
            Grouped.with(Serdes.String(), messageSerde)
        );

        // 定义时间窗口:这里用5分钟的滚动窗口(可根据需求调整)
        // 如果需要滑动窗口,可以改成 TimeWindows.of(Duration.ofMinutes(5)).advanceBy(Duration.ofMinutes(1))
        WindowedKStream<String, Message> windowedStream = groupedById.windowedBy(
            TimeWindows.of(Duration.ofMinutes(5))
                .grace(Duration.ofMinutes(10)) // 设置宽限期,处理迟到的消息
        );

        // 窗口内计数,得到KTable(键是带窗口信息的ID,值是计数)
        KTable<Windowed<String>, Long> idWindowCount = windowedStream.count(
            Materialized.as("id-window-count-store") // 指定状态存储名称,用于持久化窗口计数
        );

        // 可选:将结果输出到输出主题(把窗口信息和计数拼接成可读的字符串)
        idWindowCount.toStream()
            .map((windowedKey, count) -> KeyValue.pair(
                windowedKey.key(),
                String.format("ID: %s | 窗口开始时间: %s | 消息计数: %d",
                    windowedKey.key(),
                    windowedKey.window().startTime(),
                    count)
            ))
            .to("your-output-topic", Produced.with(Serdes.String(), Serdes.String()));

        // 启动Kafka Streams应用
        KafkaStreams streams = new KafkaStreams(builder.build(), streamsConfig);
        streams.start();

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

    // 替换成你的StreamsConfig配置方法
    private static StreamsConfig getStreamsConfig() {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "id-window-count-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());
        return new StreamsConfig(props);
    }
}
关键细节提醒
  • 事件时间vs处理时间:如果你的time字段是消息的实际产生时间(事件时间),而不是Kafka Streams处理消息的时间,你需要自定义时间提取器,从time字段解析出时间戳:
    import org.apache.kafka.clients.consumer.ConsumerRecord;
    import org.apache.kafka.streams.processor.TimestampExtractor;
    import java.time.LocalTime;
    import java.time.format.DateTimeFormatter;
    
    public class CustomTimestampExtractor implements TimestampExtractor {
        private static final DateTimeFormatter TIME_FORMATTER = DateTimeFormatter.ofPattern("HH:mm");
    
        @Override
        public long extract(ConsumerRecord<Object, Object> record, long partitionTime) {
            Message message = (Message) record.value();
            // 注意:这里仅解析了HH:MM,没有日期,实际使用时建议消息携带完整时间戳(比如ISO格式)
            LocalTime localTime = LocalTime.parse(message.getTime(), TIME_FORMATTER);
            // 结合当前日期转换为时间戳(仅示例,跨天场景会有问题)
            return localTime.atDate(java.time.LocalDate.now()).toInstant(java.time.ZoneOffset.UTC).toEpochMilli();
        }
    }
    
    然后在StreamsConfig中配置:
    props.put(StreamsConfig.DEFAULT_TIMESTAMP_EXTRACTOR_CLASS_CONFIG, CustomTimestampExtractor.class);
    
  • 窗口类型选择:滚动窗口(每个窗口不重叠)适合统计固定时间段的独立计数;滑动窗口(窗口向前滑动)适合连续的时间段统计,比如每1分钟统计过去5分钟的消息数。
  • 状态存储:Materialized.as指定的状态存储会持久化窗口计数,即使应用重启也能恢复状态,Kafka Streams会自动清理过期窗口的状态。

希望这些内容能帮到你,如果还有细节需要调整,可以根据自己的业务需求修改窗口大小、时间提取逻辑或者输出格式哦!

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:23:16