Код: Выделить всё
KStream inputStream = streamsBuilder.stream("kafka-topic", Consumed.with(Serdes.String(), Serdes.String()));
Materialized with = Materialized.with(Serdes.String(), STRING_LIST_SERDE);
KStream outputStream = inputStream
.groupByKey()
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofSeconds(2)))
.aggregate(
ArrayList::new,
(key, string, aggregate) -> {
aggregate.add(string);
return aggregate;
}, with)
.toStream();
Код: Выделить всё
outputStreamТеперь, кроме того, я хочу агрегировать сообщения до определенного предела, скажем, до тех пор, пока список не будет не превышает 50 по размеру.
Если список стал больше 50 во время агрегирования, я хочу каким-то образом разделить его на дополнительный список.
По сути, результат I Я хочу получить массив сообщений до предела по размеру (например, 50) и до определенного периода времени, в зависимости от того, что наступит раньше.
Чего мне здесь не хватает? чтобы этого добиться?
Подробнее здесь: https://stackoverflow.com/questions/766 ... -timeframe