Я изучаю Apache Kafka Stream SPI. Мне интересно, есть ли способ выполнить асинхронный код внутри метода mapValues. Например, для получения данных из внешнего хранилища. Есть ли способ взаимодействия с Kaska Streams в реактивном стиле цикла событий?
StreamsBuilder streamsBuilder = new StreamsBuilder();
streamsBuilder
.stream("SOURCE_TOPIC", Consumed.with(Serdes.String(), Serdes.String()))
.mapValues((readOnlyKey, value) -> value.toUpperCase())
.to("DESTINATION_TOPIC", Produced.with(Serdes.String(), Serdes.String()));
var topology = streamsBuilder.build()
var kafkaStreams = new KafkaStreams(topology, properties);
как заменить этот код MapValues
value.toUpperCase()
с:
CompletableFuture.completedFuture(value.toUpperCase())
Подробнее здесь: https://stackoverflow.com/questions/791 ... h-java-api