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

从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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 09:55:23