Kafka Streams事件加窗统计5秒窗口cat词频编译错误如何解决
错误原因
- 类型声明错误:调用
windowedBy()方法后,生成的KTable的key类型是Windowed<String>,不是你声明的String类型,直接触发类型不匹配编译错误 - Filter逻辑错误:你当前写的filter判断
value.equals("cat")完全不符合逻辑,这里的value是count统计出来的Long类型计数,你要过滤的是分组的单词等于cat,需要取Windowed类型key里的实际key值判断 - 正则转义错误:Java字符串中正则的
\W需要双反斜杠转义,原代码的"\W+"属于非法写法,运行时也会报错 - 输出序列化不匹配:如果直接输出Windowed类型的key,你配置的
Serdes.String()序列化规则无法匹配,需要提前把Windowed的key转成普通字符串
修正后代码
只需要修改createWordCountStream方法即可:
static void createWordCountStream(final StreamsBuilder builder) { final KStream<String, String> source = builder.stream(INPUT_TOPIC); // 修正counts变量的泛型类型,匹配windowedBy之后的返回类型 final KTable<Windowed<String>, Long> counts = source.flatMapValues(value -> Arrays.asList(value.toLowerCase(Locale.getDefault()).split("\\W+"))) .groupBy((key, value) -> value) .windowedBy(TimeWindows.of(Duration.ofSeconds(5))) .count() // 修正过滤逻辑,取Windowed包装的实际单词判断是否为cat .filter((windowedKey, count) -> windowedKey.key().equals("cat")); counts.toStream() // 提取Windowed中的实际单词作为输出key,匹配String序列化规则 .selectKey((windowedKey, count) -> windowedKey.key()) // 如果需要输出value直接是"cat xxx"的格式,打开下面这行注释,同时把to方法里的Serdes.Long()改成Serdes.String() // .mapValues(count -> "cat " + count) .to(OUTPUT_TOPIC, Produced.with(Serdes.String(), Serdes.Long())); }
内容的提问来源于stack exchange,提问作者Saranchai Angkawinijwong
相关产品推荐
相关产品推荐

