基于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
相关产品推荐
相关产品推荐

