kafka streams в runtime менять входящий/исходящий топик
Как можно в процессе работы приложения добавлять входящий топик и изменять исходящий топик? В зависимости от того с каким входящим топиком в данный момент идет работа, должен меняться исходящий топик.
final Serde<byte[]> byteArraySerde = Serdes.ByteArray();
final Serde<String> stringSerde = Serdes.String();
final StreamsBuilder builder = new StreamsBuilder();
final KStream<byte[], String> textLines = builder
.stream(prop.getProperty("kafka.topic.in"), Consumed.with(byteArraySerde, stringSerde));
final KStream<byte[], String> processed = textLines
.filter(MetaModelProcessor.filter())
.mapValues(MetaModelProcessor.getMetaModel());
processed.to(prop.getProperty("kafka.topic.out"));
final org.apache.kafka.streams.KafkaStreams streams = new org.apache.kafka.streams.KafkaStreams(builder.build(), new KafkaStreamsConfig(prop.getProperty("kafka.app.id.config"), prop.getProperty("kafka.client.id.config"), prop.getProperty("kafka.server")).getStreamsConfiguration());
streams.cleanUp();
streams.start();
Runtime.getRuntime().addShutdownHook(new Thread(streams::close));