KafkaStream is not closing correctly












0















I'm closing a KafkaStream when I need based on a certain condition:



Closing:



if(kafkaStream == null) {
logger.info("KafkaStream already closed?");
} else {
boolean closed = kafkaStream.close(10L, TimeUnit.SECONDS);
if(closed) {
kafkaStream = null;
logger.info("KafkaStream closed");
} else {
logger.info("KafkaStream could not closed");
}
}


Starting:



if(kafkaStream == null) {
logger.info("KafkaStream is starting");
kafkaStream = KafkaManager.getInstance().getStream(this.getConfigFilePath(),
this,
this.getTopic()
);
kafkaStream.start();
logger.info("KafkaStream is started");
}


In my implementation of Processor, process(String key, byte value) is still called although successfully closing stream:



public abstract class BaseKafkaProcessor implements Processor<String, byte> {
private static Logger logger = LogManager.getLogger(BaseKafkaProcessor.class);
private ProcessorContext context;


private ProcessorContext getContext() {
return context;
}

@Override
public void init(ProcessorContext context) {
this.context = context;
this.context.schedule(1000);
}


@Override
public void process(String key, byte value) {
try {
String topic = key.split("-")[0];
byte uncompressed = GzipCompressionUtil.uncompress(value);
String json = new String(uncompressed, "UTF-8");
processRecord(topic, json);
this.getContext().commit();
} catch (Exception e) {
logger.error("Error processing json", e);
}
}

protected abstract void processRecord(String topic, String json);

@Override
public void punctuate(long timestamp) {
this.getContext().commit();
}

@Override
public void close() {
this.getContext().commit();
}
}


My configuration for KafkaStreams:



application.id=dv_ws_in_app_activity_dev4
bootstrap.servers=VLXH1
auto.offset.reset=latest
num.stream.threads=1
key.serde=org.apache.kafka.common.serialization.Serdes$StringSerde
value.serde=org.apache.kafka.common.serialization.Serdes$ByteArraySerde
poll.ms = 100
commit.interval.ms=1000
state.dir=../../temp/kafka-state-dir


This client application uses the version 0.11.0.1 of Kafka.










share|improve this question



























    0















    I'm closing a KafkaStream when I need based on a certain condition:



    Closing:



    if(kafkaStream == null) {
    logger.info("KafkaStream already closed?");
    } else {
    boolean closed = kafkaStream.close(10L, TimeUnit.SECONDS);
    if(closed) {
    kafkaStream = null;
    logger.info("KafkaStream closed");
    } else {
    logger.info("KafkaStream could not closed");
    }
    }


    Starting:



    if(kafkaStream == null) {
    logger.info("KafkaStream is starting");
    kafkaStream = KafkaManager.getInstance().getStream(this.getConfigFilePath(),
    this,
    this.getTopic()
    );
    kafkaStream.start();
    logger.info("KafkaStream is started");
    }


    In my implementation of Processor, process(String key, byte value) is still called although successfully closing stream:



    public abstract class BaseKafkaProcessor implements Processor<String, byte> {
    private static Logger logger = LogManager.getLogger(BaseKafkaProcessor.class);
    private ProcessorContext context;


    private ProcessorContext getContext() {
    return context;
    }

    @Override
    public void init(ProcessorContext context) {
    this.context = context;
    this.context.schedule(1000);
    }


    @Override
    public void process(String key, byte value) {
    try {
    String topic = key.split("-")[0];
    byte uncompressed = GzipCompressionUtil.uncompress(value);
    String json = new String(uncompressed, "UTF-8");
    processRecord(topic, json);
    this.getContext().commit();
    } catch (Exception e) {
    logger.error("Error processing json", e);
    }
    }

    protected abstract void processRecord(String topic, String json);

    @Override
    public void punctuate(long timestamp) {
    this.getContext().commit();
    }

    @Override
    public void close() {
    this.getContext().commit();
    }
    }


    My configuration for KafkaStreams:



    application.id=dv_ws_in_app_activity_dev4
    bootstrap.servers=VLXH1
    auto.offset.reset=latest
    num.stream.threads=1
    key.serde=org.apache.kafka.common.serialization.Serdes$StringSerde
    value.serde=org.apache.kafka.common.serialization.Serdes$ByteArraySerde
    poll.ms = 100
    commit.interval.ms=1000
    state.dir=../../temp/kafka-state-dir


    This client application uses the version 0.11.0.1 of Kafka.










    share|improve this question

























      0












      0








      0








      I'm closing a KafkaStream when I need based on a certain condition:



      Closing:



      if(kafkaStream == null) {
      logger.info("KafkaStream already closed?");
      } else {
      boolean closed = kafkaStream.close(10L, TimeUnit.SECONDS);
      if(closed) {
      kafkaStream = null;
      logger.info("KafkaStream closed");
      } else {
      logger.info("KafkaStream could not closed");
      }
      }


      Starting:



      if(kafkaStream == null) {
      logger.info("KafkaStream is starting");
      kafkaStream = KafkaManager.getInstance().getStream(this.getConfigFilePath(),
      this,
      this.getTopic()
      );
      kafkaStream.start();
      logger.info("KafkaStream is started");
      }


      In my implementation of Processor, process(String key, byte value) is still called although successfully closing stream:



      public abstract class BaseKafkaProcessor implements Processor<String, byte> {
      private static Logger logger = LogManager.getLogger(BaseKafkaProcessor.class);
      private ProcessorContext context;


      private ProcessorContext getContext() {
      return context;
      }

      @Override
      public void init(ProcessorContext context) {
      this.context = context;
      this.context.schedule(1000);
      }


      @Override
      public void process(String key, byte value) {
      try {
      String topic = key.split("-")[0];
      byte uncompressed = GzipCompressionUtil.uncompress(value);
      String json = new String(uncompressed, "UTF-8");
      processRecord(topic, json);
      this.getContext().commit();
      } catch (Exception e) {
      logger.error("Error processing json", e);
      }
      }

      protected abstract void processRecord(String topic, String json);

      @Override
      public void punctuate(long timestamp) {
      this.getContext().commit();
      }

      @Override
      public void close() {
      this.getContext().commit();
      }
      }


      My configuration for KafkaStreams:



      application.id=dv_ws_in_app_activity_dev4
      bootstrap.servers=VLXH1
      auto.offset.reset=latest
      num.stream.threads=1
      key.serde=org.apache.kafka.common.serialization.Serdes$StringSerde
      value.serde=org.apache.kafka.common.serialization.Serdes$ByteArraySerde
      poll.ms = 100
      commit.interval.ms=1000
      state.dir=../../temp/kafka-state-dir


      This client application uses the version 0.11.0.1 of Kafka.










      share|improve this question














      I'm closing a KafkaStream when I need based on a certain condition:



      Closing:



      if(kafkaStream == null) {
      logger.info("KafkaStream already closed?");
      } else {
      boolean closed = kafkaStream.close(10L, TimeUnit.SECONDS);
      if(closed) {
      kafkaStream = null;
      logger.info("KafkaStream closed");
      } else {
      logger.info("KafkaStream could not closed");
      }
      }


      Starting:



      if(kafkaStream == null) {
      logger.info("KafkaStream is starting");
      kafkaStream = KafkaManager.getInstance().getStream(this.getConfigFilePath(),
      this,
      this.getTopic()
      );
      kafkaStream.start();
      logger.info("KafkaStream is started");
      }


      In my implementation of Processor, process(String key, byte value) is still called although successfully closing stream:



      public abstract class BaseKafkaProcessor implements Processor<String, byte> {
      private static Logger logger = LogManager.getLogger(BaseKafkaProcessor.class);
      private ProcessorContext context;


      private ProcessorContext getContext() {
      return context;
      }

      @Override
      public void init(ProcessorContext context) {
      this.context = context;
      this.context.schedule(1000);
      }


      @Override
      public void process(String key, byte value) {
      try {
      String topic = key.split("-")[0];
      byte uncompressed = GzipCompressionUtil.uncompress(value);
      String json = new String(uncompressed, "UTF-8");
      processRecord(topic, json);
      this.getContext().commit();
      } catch (Exception e) {
      logger.error("Error processing json", e);
      }
      }

      protected abstract void processRecord(String topic, String json);

      @Override
      public void punctuate(long timestamp) {
      this.getContext().commit();
      }

      @Override
      public void close() {
      this.getContext().commit();
      }
      }


      My configuration for KafkaStreams:



      application.id=dv_ws_in_app_activity_dev4
      bootstrap.servers=VLXH1
      auto.offset.reset=latest
      num.stream.threads=1
      key.serde=org.apache.kafka.common.serialization.Serdes$StringSerde
      value.serde=org.apache.kafka.common.serialization.Serdes$ByteArraySerde
      poll.ms = 100
      commit.interval.ms=1000
      state.dir=../../temp/kafka-state-dir


      This client application uses the version 0.11.0.1 of Kafka.







      java apache-kafka apache-kafka-streams






      share|improve this question













      share|improve this question











      share|improve this question




      share|improve this question










      asked Nov 14 '18 at 8:26









      OzgurOzgur

      2,829625




      2,829625
























          0






          active

          oldest

          votes











          Your Answer






          StackExchange.ifUsing("editor", function () {
          StackExchange.using("externalEditor", function () {
          StackExchange.using("snippets", function () {
          StackExchange.snippets.init();
          });
          });
          }, "code-snippets");

          StackExchange.ready(function() {
          var channelOptions = {
          tags: "".split(" "),
          id: "1"
          };
          initTagRenderer("".split(" "), "".split(" "), channelOptions);

          StackExchange.using("externalEditor", function() {
          // Have to fire editor after snippets, if snippets enabled
          if (StackExchange.settings.snippets.snippetsEnabled) {
          StackExchange.using("snippets", function() {
          createEditor();
          });
          }
          else {
          createEditor();
          }
          });

          function createEditor() {
          StackExchange.prepareEditor({
          heartbeatType: 'answer',
          autoActivateHeartbeat: false,
          convertImagesToLinks: true,
          noModals: true,
          showLowRepImageUploadWarning: true,
          reputationToPostImages: 10,
          bindNavPrevention: true,
          postfix: "",
          imageUploader: {
          brandingHtml: "Powered by u003ca class="icon-imgur-white" href="https://imgur.com/"u003eu003c/au003e",
          contentPolicyHtml: "User contributions licensed under u003ca href="https://creativecommons.org/licenses/by-sa/3.0/"u003ecc by-sa 3.0 with attribution requiredu003c/au003e u003ca href="https://stackoverflow.com/legal/content-policy"u003e(content policy)u003c/au003e",
          allowUrls: true
          },
          onDemand: true,
          discardSelector: ".discard-answer"
          ,immediatelyShowMarkdownHelp:true
          });


          }
          });














          draft saved

          draft discarded


















          StackExchange.ready(
          function () {
          StackExchange.openid.initPostLogin('.new-post-login', 'https%3a%2f%2fstackoverflow.com%2fquestions%2f53295821%2fkafkastream-is-not-closing-correctly%23new-answer', 'question_page');
          }
          );

          Post as a guest















          Required, but never shown

























          0






          active

          oldest

          votes








          0






          active

          oldest

          votes









          active

          oldest

          votes






          active

          oldest

          votes
















          draft saved

          draft discarded




















































          Thanks for contributing an answer to Stack Overflow!


          • Please be sure to answer the question. Provide details and share your research!

          But avoid



          • Asking for help, clarification, or responding to other answers.

          • Making statements based on opinion; back them up with references or personal experience.


          To learn more, see our tips on writing great answers.




          draft saved


          draft discarded














          StackExchange.ready(
          function () {
          StackExchange.openid.initPostLogin('.new-post-login', 'https%3a%2f%2fstackoverflow.com%2fquestions%2f53295821%2fkafkastream-is-not-closing-correctly%23new-answer', 'question_page');
          }
          );

          Post as a guest















          Required, but never shown





















































          Required, but never shown














          Required, but never shown












          Required, but never shown







          Required, but never shown

































          Required, but never shown














          Required, but never shown












          Required, but never shown







          Required, but never shown







          Popular posts from this blog

          Xamarin.iOS Cant Deploy on Iphone

          Glorious Revolution

          Dulmage-Mendelsohn matrix decomposition in Python