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

基于Kafka Streams的小时级消息计数Java示例需求

Kafka Streams 每小时消息计数示例应用

核心实现思路

使用Kafka Streams的1小时滚动窗口对输入主题的所有消息进行全局计数,通过suppress操作确保仅在窗口结束后输出最终计数,最后在窗口结果回调中调用方法生成包含计数的JSON payload。

依赖配置(Maven)

确保引入Kafka Streams和Jackson依赖用于JSON处理:

<dependencies>
    <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka-streams</artifactId>
        <version>3.6.1</version> <!-- 使用最新稳定版本 -->
    </dependency>
    <dependency>
        <groupId>com.fasterxml.jackson.core</groupId>
        <artifactId>jackson-databind</artifactId>
        <version>2.15.2</version>
    </dependency>
</dependencies>

完整Java代码示例

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
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.*;
import org.apache.kafka.streams.kstream.Suppressed.BufferConfig;

import java.time.Duration;
import java.util.HashMap;
import java.util.Map;
import java.util.Properties;

public class HourlyMessageCountApp {

    public static void main(String[] args) {
        // 配置Kafka Streams参数
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "hourly-message-count-app");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); // 替换为你的Kafka地址
        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();

        // 读取输入主题(替换为你的目标主题名)
        KStream<String, String> inputStream = builder.stream("input-topic");

        // 全局计数:将所有消息映射到同一个键,实现全主题统计
        KTable<Windowed<String>, Long> hourlyCountTable = inputStream
                .selectKey((key, value) -> "global-total") // 统一键值,用于全局聚合
                .groupByKey()
                .windowedBy(TimeWindows.of(Duration.ofHours(1))
                        .grace(Duration.ofMinutes(5))) // 允许5分钟延迟处理迟到消息
                .count();

        // 抑制中间输出,仅在窗口关闭后发送最终计数
        KTable<Windowed<String>, Long> finalCountTable = hourlyCountTable
                .suppress(Suppressed.untilWindowCloses(BufferConfig.unbounded()));

        // 处理窗口结束事件,生成目标payload
        finalCountTable.toStream().foreach((windowedKey, count) -> {
            try {
                String payload = generateCountPayload(windowedKey, count);
                // 这里可根据需求将payload发送到输出主题、存储或调用外部服务
                System.out.println("窗口结束生成Payload: " + payload);
            } catch (JsonProcessingException e) {
                e.printStackTrace();
            }
        });

        // 启动流应用
        KafkaStreams streams = new KafkaStreams(builder.build(), props);
        streams.start();

        // 注册JVM关闭钩子,优雅停止应用
        Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
    }

    /**
     * 根据窗口信息和计数生成JSON格式的payload
     */
    private static String generateCountPayload(Windowed<String> windowedKey, Long count) throws JsonProcessingException {
        ObjectMapper objectMapper = new ObjectMapper();
        Map<String, Object> payload = new HashMap<>();
        payload.put("window_start_timestamp", windowedKey.window().startTime().toEpochMilli());
        payload.put("window_end_timestamp", windowedKey.window().endTime().toEpochMilli());
        payload.put("window_human_start", windowedKey.window().startTime().toString());
        payload.put("window_human_end", windowedKey.window().endTime().toString());
        payload.put("total_message_count", count);
        payload.put("payload_version", "1.0");
        return objectMapper.writeValueAsString(payload);
    }
}

关键细节说明

  • 滚动窗口配置:TimeWindows.of(Duration.ofHours(1))创建无重叠的1小时窗口,grace(Duration.ofMinutes(5))允许迟到5分钟的消息被计入对应窗口,避免网络延迟导致计数偏差。
  • 全局计数实现:通过selectKey将所有消息映射到同一个固定键,实现对整个主题的消息计数,而非按原始key分组。
  • suppress操作:确保仅在窗口完全关闭后输出最终计数,而非实时输出中间累加值,符合“每小时结束时生成结果”的需求。
  • Payload生成:generateCountPayload方法根据窗口时间范围和计数生成结构化JSON,可根据业务需求调整字段内容。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 04:07:37