Всем привет!
Возник вопрос почему не работает окно и reduce в kafka streams, может каких настроек не хватает.
Проблема в следующем, я ожидаю что окно соберет все сообщения и сгруппирует по нужному мне полю из объекта, а reduce уже отбросит лишнее и в foreach попадут только нужные записи, но почему то в foreach летит все что прилетело в топик.
Вот пример
input
.peek { key, value ->
logger.debug { }
}
.groupBy(
{ _, value -> value.requestId.toString() },
withSerializer
)
.windowedBy(TimeWindows.of(Duration.ofMinute(1)).grace(
Duration.ZERO))
.reduce(
{ first, second ->
if (first > second) {
return@reduce first
} else {
return@reduce second
}
}, withMaterializer
)
.toStream()
.foreach {}