Как добавить пользовательский StateStore к процессору Kafka Streams DSL?

Для одного из моих потоковых приложений Kafka мне нужно использовать функции как DSL, так и Processor API. Мой поток потокового приложения

source -> selectKey -> filter -> aggregate (on a window) -> sink

После агрегации мне нужно отправить ЕДИНОЕ агрегированное сообщение в приемник. Поэтому я определяю свою топологию, как показано ниже

KStreamBuilder builder = new KStreamBuilder();
KStream<String, String> source = builder.stream(source_stream);
source.selectKey(new MyKeyValueMapper())
      .filterNot((k,v) -> k.equals("UnknownGroup"))
      .process(() -> new MyProcessor());

Я определяю обычайStateStore и зарегистрируйте его на моем процессоре, как показано ниже

public class MyProcessor implements Processor<String, String> {

    private ProcessorContext context = null;
    Serde<HashMapStore> invSerde = Serdes.serdeFrom(invJsonSerializer, invJsonDeserializer);


    KeyValueStore<String, HashMapStore> invStore = (KeyValueStore) Stores.create("invStore")
        .withKeys(Serdes.String())
        .withValues(invSerde)
        .persistent()
        .build()
        .get();

    public MyProcessor() {
    }

    @Override
    public void init(ProcessorContext context) {
        this.context = context;
        this.context.register(invStore, false, null); // register the store
        this.context.schedule(10 * 60 * 1000L);
    }

    @Override
    public void process(String partitionKey, String message) {
        try {
            MessageModel smb = new MessageModel(message);
            HashMapStore oldStore = invStore.get(partitionKey);
            if (oldStore == null) {
                oldStore = new HashMapStore();
            }
            oldStore.addSmb(smb);
            invStore.put(partitionKey, oldStore);
        } catch (Exception e) {
           e.printStackTrace();
        }
    }

    @Override
    public void punctuate(long timestamp) {
       // processes all the messages in the state store and sends single aggregate message
    }


    @Override
    public void close() {
        invStore.close();
    }
}

Когда я запускаю приложение, я получаюjava.lang.NullPointerException

Исключение в потоке "StreamThread-18" java.lang.NullPointerException в org.apache.kafka.streams.state.internals.MeteredKeyValueStore.flush (MeteredKeyValueStore.java:167) в org.apache.kafka.streams.processor.internalsProcessoror .flush (ProcessorStateManager.java:332) в org.apache.kafka.streams.processor.internals.StreamTask.commit (StreamTask.java:252) в org.apache.kafka.streams.processor.internals.StreamThread.commitOne (StreamThread .java: 446) в org.apache.kafka.streams.processor.internals.StreamThread.commitAll (StreamThread.java:434) в org.apache.kafka.streams.processor.internals.StreamThread.maybeCommit (StreamThread.java:422). ) в org.apache.kafka.streams.processor.internals.StreamThread.runLoop (StreamThread.java:340) в org.apache.kafka.streams.processor.internals.StreamThread.run (StreamThread.java:218)

Есть идеи, что здесь происходит?

Ответы на вопрос(1)

Ваш ответ на вопрос