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));

Ответы (0 шт):