【发布时间】:2020-10-26 19:04:21
【问题描述】:
我有其中包含用户 ID 的流式事件。我想计算在特定时间内有多少不同的用户生成事件。但是,我是 Kafka 的初学者,我无法应对这个问题。
1 分钟内的示例事件;
{"event_name": "viewProduct", "user_id": "12"}
{"event_name": "viewProductDetails", "user_id": "23"}
{"event_name": "viewProductComments", "user_id": "12"}
{"event_name": "viewProduct", "user_id": "23"}
{"event_name": "viewProductComments", "user_id": "32"}
根据上述事件,我的代码应该生成有 3 个活跃用户。
我的方法如下,但是这个解决方案不能消除来自同一用户的多个事件并多次计算同一用户。
builder.stream("orders") // read from orders toic
.mapValues(v -> { // get user_id via json parser
JsonNode jsonNode = null;
try {
jsonNode = objectMapper.readTree((String) v);
return jsonNode.get("user_id").asText();
} catch (JsonProcessingException e) {
e.printStackTrace();
}
return "";
})
.selectKey((k, v) -> "1") // put same key to every user_id
.groupByKey()
.windowedBy(TimeWindows.of(Duration.ofSeconds(1))) // use time windows
.count() // count values
【问题讨论】:
-
为什么要在所有记录上设置相同的键(“1”)?
-
我想将所有数据分组在同一个键下并轻松计数。但是,我知道我的方法不好并且没有弄清楚问题。
-
顺便说一下,如果你使用 JSONSerde 作为值反序列化器,你就不需要 ObjectMapper
标签: apache-kafka apache-kafka-streams