从KSQLDB迁移至Java KStreams:多字段分组统计访问次数拓扑实现
Kafka Streams 多字段分组统计实现方案
要实现按user_id、object_id、日期(从entrytimestamp提取)分组统计用户每日查看特定对象的次数,核心是构建复合分组键或使用滚动窗口,结合聚合操作维护状态计数。以下是完整实现步骤及代码示例:
1. 定义数据模型
创建与输入JSON结构匹配的POJO类,用于反序列化Kafka消息:
import com.fasterxml.jackson.annotation.JsonProperty; import java.time.LocalDateTime; import java.time.format.DateTimeFormatter; public class ViewEvent { @JsonProperty("entrytimestamp") private String entryTimestamp; @JsonProperty("user_id") private String userId; @JsonProperty("object_id") private String objectId; // Jackson反序列化需要无参构造器 public ViewEvent() {} // Getter & Setter public String getEntryTimestamp() { return entryTimestamp; } public void setEntryTimestamp(String entryTimestamp) { this.entryTimestamp = entryTimestamp; } public String getUserId() { return userId; } public void setUserId(String userId) { this.userId = userId; } public String getObjectId() { return objectId; } public void setObjectId(String objectId) { this.objectId = objectId; } // 提取日期(格式:yyyy-MM-dd) public String getDate() { DateTimeFormatter formatter = DateTimeFormatter.ISO_LOCAL_DATE_TIME; LocalDateTime dateTime = LocalDateTime.parse(entryTimestamp, formatter); return dateTime.toLocalDate().toString(); } }
2. 自定义复合键类
创建包含user_id、object_id、日期的复合键类,用于分组:
import java.io.Serializable; import java.util.Objects; public class UserObjectDateKey implements Serializable { private String userId; private String objectId; private String date; public UserObjectDateKey() {} public UserObjectDateKey(String userId, String objectId, String date) { this.userId = userId; this.objectId = objectId; this.date = date; } // Getter & Setter public String getUserId() { return userId; } public void setUserId(String userId) { this.userId = userId; } public String getObjectId() { return objectId; } public void setObjectId(String objectId) { this.objectId = objectId; } public String getDate() { return date; } public void setDate(String date) { this.date = date; } // 必须重写equals和hashCode,确保分组逻辑正确 @Override public boolean equals(Object o) { if (this == o) return true; if (o == null || getClass() != o.getClass()) return false; UserObjectDateKey that = (UserObjectDateKey) o; return Objects.equals(userId, that.userId) && Objects.equals(objectId, that.objectId) && Objects.equals(date, that.date); } @Override public int hashCode() { return Objects.hash(userId, objectId, date); } }
3. 自定义JSON Serde
Kafka Streams默认不支持自定义类的序列化,因此实现通用JSON Serde:
import org.apache.kafka.common.serialization.Deserializer; import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.common.serialization.Serializer; import com.fasterxml.jackson.databind.ObjectMapper; public class JsonSerde<T> implements Serde<T> { private static final ObjectMapper objectMapper = new ObjectMapper(); private final Class<T> type; public JsonSerde(Class<T> type) { this.type = type; } @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); } }; } }
4. 构建Stream拓扑
方案一:复合键直接分组
将日期作为复合键的一部分,直接按复合键分组计数:
import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.Topology; import org.apache.kafka.streams.kstream.Consumed; import org.apache.kafka.streams.kstream.Grouped; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.Materialized; import org.apache.kafka.streams.kstream.Produced; import java.util.Properties; public class ViewCountStream { public static void main(String[] args) { Properties props = new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "user-object-daily-view-count"); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); StreamsBuilder builder = new StreamsBuilder(); // 从源Topic读取数据 KStream<String, ViewEvent> sourceStream = builder.stream( "view-events-topic", Consumed.with(new JsonSerde<>(String.class), new JsonSerde<>(ViewEvent.class)) ); // 转换为复合键+计数增量(1)的流 KStream<UserObjectDateKey, Long> countStream = sourceStream .map((key, event) -> { UserObjectDateKey compositeKey = new UserObjectDateKey( event.getUserId(), event.getObjectId(), event.getDate() ); return new org.apache.kafka.streams.KeyValue<>(compositeKey, 1L); }); // 分组聚合并输出结果 countStream .groupByKey(Grouped.with(new JsonSerde<>(UserObjectDateKey.class), new JsonSerde<>(Long.class))) .count(Materialized.as("user-object-daily-view-count-store")) // 指定状态存储名称 .toStream() .to("daily-view-counts-topic", Produced.with( new JsonSerde<>(UserObjectDateKey.class), new JsonSerde<>(Long.class) )); Topology topology = builder.build(); KafkaStreams streams = new KafkaStreams(topology, props); streams.start(); Runtime.getRuntime().addShutdownHook(new Thread(streams::close)); } }
方案二:滚动窗口分组
若更倾向于使用窗口语义,可按user_id+object_id分组,结合1天滚动窗口计数:
// 先定义仅包含user_id和object_id的键类 class UserObjectKey implements Serializable { private String userId; private String objectId; // 构造器、Getter&Setter、equals&hashCode 实现同UserObjectDateKey } // 拓扑中替换为窗口逻辑 sourceStream .map((key, event) -> new org.apache.kafka.streams.KeyValue<>( new UserObjectKey(event.getUserId(), event.getObjectId()), 1L )) .groupByKey(Grouped.with(new JsonSerde<>(UserObjectKey.class), new JsonSerde<>(Long.class))) .windowedBy(TimeWindows.of(Duration.ofDays(1))) // 1天滚动窗口 .count(Materialized.as("user-object-windowed-view-count-store")) .toStream() .map((windowedKey, count) -> { // 从窗口起始时间提取日期 String date = windowedKey.window().startTime().toLocalDateTime().toLocalDate().toString(); UserObjectDateKey resultKey = new UserObjectDateKey( windowedKey.key().getUserId(), windowedKey.key().getObjectId(), date ); return new org.apache.kafka.streams.KeyValue<>(resultKey, count); }) .to("daily-view-counts-topic", Produced.with( new JsonSerde<>(UserObjectDateKey.class), new JsonSerde<>(Long.class) ));
关键注意事项
- 时区处理:若
entrytimestamp包含时区,需改用ZonedDateTime解析,避免日期计算错误。 - 事件时间配置:生产环境建议使用事件时间而非默认的处理时间,需自定义
TimestampExtractor从entrytimestamp提取事件时间并配置到StreamsConfig。 - 状态存储:指定状态存储名称后,Kafka Streams会自动维护状态,生产环境需配置持久化状态目录或远程状态存储。
- Serde正确性:自定义类的Serde必须正确实现,否则会导致序列化错误,进而分组失败。
内容的提问来源于stack exchange,提问作者user3008410
相关产品推荐
相关产品推荐

