Proper way to programmatically stop an Alpakka Kafka stream2019 Community Moderator ElectionData Modeling with Kafka? Topics and PartitionsCan I run Kafka Streams Application on the same machine as of Kafka Broker?Akka Streams Reactive Kafka - OutOfMemoryError under high loadAkka Kafka stream supervison strategy not workingAkka Streams KillSwitch in alpakka jmsGracefully restart a Reactive-Kafka Consumer Stream on failureKafka Streams: Kafka Streams application stuck rebalancingAlpakka/Kafka - Partitions consumed faster than othersActorSystem shutdown in akka streamHow to control/pause akka streams flow in sub source/stream

Loading the leaflet Map in Lightning Web Component

What exactly term 'companion plants' means?

What (if any) is the reason to buy in small local stores?

Asserting that Atheism and Theism are both faith based positions

Recruiter wants very extensive technical details about all of my previous work

How does 取材で訪れた integrate into this sentence?

What does "mu" mean as an interjection?

Help rendering a complicated sum/product formula

Knife as defense against stray dogs

Inhabiting Mars versus going straight for a Dyson swarm

Synchronized implementation of a bank account in Java

How can an organ that provides biological immortality be unable to regenerate?

두음법칙 - When did North and South diverge in pronunciation of initial ㄹ?

What does "Four-F." mean?

How are passwords stolen from companies if they only store hashes?

Is it insecure to send a password in a `curl` command?

Can a medieval gyroplane be built?

Should I be concerned about student access to a test bank?

How to generate binary array whose elements with values 1 are randomly drawn

Variable completely messes up echoed string

Fewest number of steps to reach 200 using special calculator

Generic TVP tradeoffs?

In Aliens, how many people were on LV-426 before the Marines arrived​?

Is honey really a supersaturated solution? Does heating to un-crystalize redissolve it or melt it?



Proper way to programmatically stop an Alpakka Kafka stream



2019 Community Moderator ElectionData Modeling with Kafka? Topics and PartitionsCan I run Kafka Streams Application on the same machine as of Kafka Broker?Akka Streams Reactive Kafka - OutOfMemoryError under high loadAkka Kafka stream supervison strategy not workingAkka Streams KillSwitch in alpakka jmsGracefully restart a Reactive-Kafka Consumer Stream on failureKafka Streams: Kafka Streams application stuck rebalancingAlpakka/Kafka - Partitions consumed faster than othersActorSystem shutdown in akka streamHow to control/pause akka streams flow in sub source/stream










4















We are trying to use Akka Streams with Alpakka Kafka to consume a stream of events in a service. For handling event processing errors we are using Kafka autocommit and more than one queue. For example, if we have the topic user_created, which we want to consume from a products service, we also create user_created_for_products_failed and user_created_for_products_dead_letter. These two extra topics are coupled to a specific Kafka consumer group. If an event fails to be processed, it goes to the failed queue, where we try to consume again in five minutes--if it fails again it goes to dead letters.



On deployment we want to ensure that we don't lose events. So we are trying to stop the stream before stopping the application. As I said, we are using autocommit, but all of these events that are "flying" are not processed yet. Once the stream and application are stopped, we can deploy the new code and start the application again.



After reading the documentation, we have seen the KillSwitch feature. The problem that we are seeing in it is that the shutdown method returns Unit instead Future[Unit] as we expect. We are not sure that we won't lose events using it, because in tests it looks like it goes too fast to be working properly.



As a workaround, we create an ActorSystem for each stream and use the terminate method (which returns a Future[Terminate]). The problem with this solution is that we don't think that creating an ActorSystem per stream will scale well, and terminate takes a lot of time to resolve (in tests it takes up to one minute to shut down).



Have you faced a problem like this? Is there a faster way (compared to ActorSystem.terminate) to stop a stream and ensure that all the events that the Source has emitted have been processed?










share|improve this question




























    4















    We are trying to use Akka Streams with Alpakka Kafka to consume a stream of events in a service. For handling event processing errors we are using Kafka autocommit and more than one queue. For example, if we have the topic user_created, which we want to consume from a products service, we also create user_created_for_products_failed and user_created_for_products_dead_letter. These two extra topics are coupled to a specific Kafka consumer group. If an event fails to be processed, it goes to the failed queue, where we try to consume again in five minutes--if it fails again it goes to dead letters.



    On deployment we want to ensure that we don't lose events. So we are trying to stop the stream before stopping the application. As I said, we are using autocommit, but all of these events that are "flying" are not processed yet. Once the stream and application are stopped, we can deploy the new code and start the application again.



    After reading the documentation, we have seen the KillSwitch feature. The problem that we are seeing in it is that the shutdown method returns Unit instead Future[Unit] as we expect. We are not sure that we won't lose events using it, because in tests it looks like it goes too fast to be working properly.



    As a workaround, we create an ActorSystem for each stream and use the terminate method (which returns a Future[Terminate]). The problem with this solution is that we don't think that creating an ActorSystem per stream will scale well, and terminate takes a lot of time to resolve (in tests it takes up to one minute to shut down).



    Have you faced a problem like this? Is there a faster way (compared to ActorSystem.terminate) to stop a stream and ensure that all the events that the Source has emitted have been processed?










    share|improve this question


























      4












      4








      4








      We are trying to use Akka Streams with Alpakka Kafka to consume a stream of events in a service. For handling event processing errors we are using Kafka autocommit and more than one queue. For example, if we have the topic user_created, which we want to consume from a products service, we also create user_created_for_products_failed and user_created_for_products_dead_letter. These two extra topics are coupled to a specific Kafka consumer group. If an event fails to be processed, it goes to the failed queue, where we try to consume again in five minutes--if it fails again it goes to dead letters.



      On deployment we want to ensure that we don't lose events. So we are trying to stop the stream before stopping the application. As I said, we are using autocommit, but all of these events that are "flying" are not processed yet. Once the stream and application are stopped, we can deploy the new code and start the application again.



      After reading the documentation, we have seen the KillSwitch feature. The problem that we are seeing in it is that the shutdown method returns Unit instead Future[Unit] as we expect. We are not sure that we won't lose events using it, because in tests it looks like it goes too fast to be working properly.



      As a workaround, we create an ActorSystem for each stream and use the terminate method (which returns a Future[Terminate]). The problem with this solution is that we don't think that creating an ActorSystem per stream will scale well, and terminate takes a lot of time to resolve (in tests it takes up to one minute to shut down).



      Have you faced a problem like this? Is there a faster way (compared to ActorSystem.terminate) to stop a stream and ensure that all the events that the Source has emitted have been processed?










      share|improve this question
















      We are trying to use Akka Streams with Alpakka Kafka to consume a stream of events in a service. For handling event processing errors we are using Kafka autocommit and more than one queue. For example, if we have the topic user_created, which we want to consume from a products service, we also create user_created_for_products_failed and user_created_for_products_dead_letter. These two extra topics are coupled to a specific Kafka consumer group. If an event fails to be processed, it goes to the failed queue, where we try to consume again in five minutes--if it fails again it goes to dead letters.



      On deployment we want to ensure that we don't lose events. So we are trying to stop the stream before stopping the application. As I said, we are using autocommit, but all of these events that are "flying" are not processed yet. Once the stream and application are stopped, we can deploy the new code and start the application again.



      After reading the documentation, we have seen the KillSwitch feature. The problem that we are seeing in it is that the shutdown method returns Unit instead Future[Unit] as we expect. We are not sure that we won't lose events using it, because in tests it looks like it goes too fast to be working properly.



      As a workaround, we create an ActorSystem for each stream and use the terminate method (which returns a Future[Terminate]). The problem with this solution is that we don't think that creating an ActorSystem per stream will scale well, and terminate takes a lot of time to resolve (in tests it takes up to one minute to shut down).



      Have you faced a problem like this? Is there a faster way (compared to ActorSystem.terminate) to stop a stream and ensure that all the events that the Source has emitted have been processed?







      scala apache-kafka akka akka-stream alpakka






      share|improve this question















      share|improve this question













      share|improve this question




      share|improve this question








      edited Mar 7 at 17:35









      Jeffrey Chung

      14.3k62142




      14.3k62142










      asked Mar 7 at 14:28









      SergiGPSergiGP

      259313




      259313






















          1 Answer
          1






          active

          oldest

          votes


















          4














          From the documentation (emphasis mine):




          When using external offset storage, a call to Consumer.Control.shutdown() suffices to complete the Source, which starts the completion of the stream.




          val (consumerControl, streamComplete) =
          Consumer
          .plainSource(consumerSettings,
          Subscriptions.assignmentWithOffset(
          new TopicPartition(topic, 0) -> offset
          ))
          .via(businessFlow)
          .toMat(Sink.ignore)(Keep.both)
          .run()

          consumerControl.shutdown()


          Consumer.control.shutdown() returns a Future[Done]. From its Scaladoc description:




          Shutdown the consumer Source. It will wait for outstanding offset commit requests to finish before shutting down.




          Alternatively, if you're using offset storage in Kafka, use Consumer.Control.drainAndShutdown, which also returns a Future. Again from the documentation (which contains more information about what drainAndShutdown does under the covers):



          val drainingControl =
          Consumer
          .committableSource(consumerSettings.withStopTimeout(Duration.Zero), Subscriptions.topics(topic))
          .mapAsync(1) msg =>
          business(msg.record).map(_ => msg.committableOffset)

          .toMat(Committer.sink(committerSettings))(Keep.both)
          .mapMaterializedValue(DrainingControl.apply)
          .run()

          val streamComplete = drainingControl.drainAndShutdown()


          The Scaladoc description for drainAndShutdown:




          Stop producing messages from the Source, wait for stream completion and shut down the consumer Source so that all consumed messages reach the end of the stream. Failures in stream completion will be propagated, the source will be shut down anyway.







          share|improve this answer

























          • Cool dude! Thank you!

            – SergiGP
            Mar 8 at 10:41










          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%2f55046169%2fproper-way-to-programmatically-stop-an-alpakka-kafka-stream%23new-answer', 'question_page');

          );

          Post as a guest















          Required, but never shown

























          1 Answer
          1






          active

          oldest

          votes








          1 Answer
          1






          active

          oldest

          votes









          active

          oldest

          votes






          active

          oldest

          votes









          4














          From the documentation (emphasis mine):




          When using external offset storage, a call to Consumer.Control.shutdown() suffices to complete the Source, which starts the completion of the stream.




          val (consumerControl, streamComplete) =
          Consumer
          .plainSource(consumerSettings,
          Subscriptions.assignmentWithOffset(
          new TopicPartition(topic, 0) -> offset
          ))
          .via(businessFlow)
          .toMat(Sink.ignore)(Keep.both)
          .run()

          consumerControl.shutdown()


          Consumer.control.shutdown() returns a Future[Done]. From its Scaladoc description:




          Shutdown the consumer Source. It will wait for outstanding offset commit requests to finish before shutting down.




          Alternatively, if you're using offset storage in Kafka, use Consumer.Control.drainAndShutdown, which also returns a Future. Again from the documentation (which contains more information about what drainAndShutdown does under the covers):



          val drainingControl =
          Consumer
          .committableSource(consumerSettings.withStopTimeout(Duration.Zero), Subscriptions.topics(topic))
          .mapAsync(1) msg =>
          business(msg.record).map(_ => msg.committableOffset)

          .toMat(Committer.sink(committerSettings))(Keep.both)
          .mapMaterializedValue(DrainingControl.apply)
          .run()

          val streamComplete = drainingControl.drainAndShutdown()


          The Scaladoc description for drainAndShutdown:




          Stop producing messages from the Source, wait for stream completion and shut down the consumer Source so that all consumed messages reach the end of the stream. Failures in stream completion will be propagated, the source will be shut down anyway.







          share|improve this answer

























          • Cool dude! Thank you!

            – SergiGP
            Mar 8 at 10:41















          4














          From the documentation (emphasis mine):




          When using external offset storage, a call to Consumer.Control.shutdown() suffices to complete the Source, which starts the completion of the stream.




          val (consumerControl, streamComplete) =
          Consumer
          .plainSource(consumerSettings,
          Subscriptions.assignmentWithOffset(
          new TopicPartition(topic, 0) -> offset
          ))
          .via(businessFlow)
          .toMat(Sink.ignore)(Keep.both)
          .run()

          consumerControl.shutdown()


          Consumer.control.shutdown() returns a Future[Done]. From its Scaladoc description:




          Shutdown the consumer Source. It will wait for outstanding offset commit requests to finish before shutting down.




          Alternatively, if you're using offset storage in Kafka, use Consumer.Control.drainAndShutdown, which also returns a Future. Again from the documentation (which contains more information about what drainAndShutdown does under the covers):



          val drainingControl =
          Consumer
          .committableSource(consumerSettings.withStopTimeout(Duration.Zero), Subscriptions.topics(topic))
          .mapAsync(1) msg =>
          business(msg.record).map(_ => msg.committableOffset)

          .toMat(Committer.sink(committerSettings))(Keep.both)
          .mapMaterializedValue(DrainingControl.apply)
          .run()

          val streamComplete = drainingControl.drainAndShutdown()


          The Scaladoc description for drainAndShutdown:




          Stop producing messages from the Source, wait for stream completion and shut down the consumer Source so that all consumed messages reach the end of the stream. Failures in stream completion will be propagated, the source will be shut down anyway.







          share|improve this answer

























          • Cool dude! Thank you!

            – SergiGP
            Mar 8 at 10:41













          4












          4








          4







          From the documentation (emphasis mine):




          When using external offset storage, a call to Consumer.Control.shutdown() suffices to complete the Source, which starts the completion of the stream.




          val (consumerControl, streamComplete) =
          Consumer
          .plainSource(consumerSettings,
          Subscriptions.assignmentWithOffset(
          new TopicPartition(topic, 0) -> offset
          ))
          .via(businessFlow)
          .toMat(Sink.ignore)(Keep.both)
          .run()

          consumerControl.shutdown()


          Consumer.control.shutdown() returns a Future[Done]. From its Scaladoc description:




          Shutdown the consumer Source. It will wait for outstanding offset commit requests to finish before shutting down.




          Alternatively, if you're using offset storage in Kafka, use Consumer.Control.drainAndShutdown, which also returns a Future. Again from the documentation (which contains more information about what drainAndShutdown does under the covers):



          val drainingControl =
          Consumer
          .committableSource(consumerSettings.withStopTimeout(Duration.Zero), Subscriptions.topics(topic))
          .mapAsync(1) msg =>
          business(msg.record).map(_ => msg.committableOffset)

          .toMat(Committer.sink(committerSettings))(Keep.both)
          .mapMaterializedValue(DrainingControl.apply)
          .run()

          val streamComplete = drainingControl.drainAndShutdown()


          The Scaladoc description for drainAndShutdown:




          Stop producing messages from the Source, wait for stream completion and shut down the consumer Source so that all consumed messages reach the end of the stream. Failures in stream completion will be propagated, the source will be shut down anyway.







          share|improve this answer















          From the documentation (emphasis mine):




          When using external offset storage, a call to Consumer.Control.shutdown() suffices to complete the Source, which starts the completion of the stream.




          val (consumerControl, streamComplete) =
          Consumer
          .plainSource(consumerSettings,
          Subscriptions.assignmentWithOffset(
          new TopicPartition(topic, 0) -> offset
          ))
          .via(businessFlow)
          .toMat(Sink.ignore)(Keep.both)
          .run()

          consumerControl.shutdown()


          Consumer.control.shutdown() returns a Future[Done]. From its Scaladoc description:




          Shutdown the consumer Source. It will wait for outstanding offset commit requests to finish before shutting down.




          Alternatively, if you're using offset storage in Kafka, use Consumer.Control.drainAndShutdown, which also returns a Future. Again from the documentation (which contains more information about what drainAndShutdown does under the covers):



          val drainingControl =
          Consumer
          .committableSource(consumerSettings.withStopTimeout(Duration.Zero), Subscriptions.topics(topic))
          .mapAsync(1) msg =>
          business(msg.record).map(_ => msg.committableOffset)

          .toMat(Committer.sink(committerSettings))(Keep.both)
          .mapMaterializedValue(DrainingControl.apply)
          .run()

          val streamComplete = drainingControl.drainAndShutdown()


          The Scaladoc description for drainAndShutdown:




          Stop producing messages from the Source, wait for stream completion and shut down the consumer Source so that all consumed messages reach the end of the stream. Failures in stream completion will be propagated, the source will be shut down anyway.








          share|improve this answer














          share|improve this answer



          share|improve this answer








          edited Mar 7 at 17:46

























          answered Mar 7 at 14:57









          Jeffrey ChungJeffrey Chung

          14.3k62142




          14.3k62142












          • Cool dude! Thank you!

            – SergiGP
            Mar 8 at 10:41

















          • Cool dude! Thank you!

            – SergiGP
            Mar 8 at 10:41
















          Cool dude! Thank you!

          – SergiGP
          Mar 8 at 10:41





          Cool dude! Thank you!

          – SergiGP
          Mar 8 at 10:41



















          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%2f55046169%2fproper-way-to-programmatically-stop-an-alpakka-kafka-stream%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

          Thal And Out Agency railway station See also References External links Navigation menuOfficial Web Site of Pakistan RailwaysArchivedOfficial Web Site of Pakistan Railwayseeexpanding ite

          Understanding generators in Python2019 Community Moderator ElectionGenerator function not working pythonFor loop not executing two timesGenerators - Printing generated valuesWhy can a python generator only be used once?What exactly do generators do?What does the “yield” keyword do?What does “list comprehension” mean? How does it work and how can I use it?sklearn Kfold acces single fold instead of for loopIs a generator the callable? Which is the generator?Apply Border To Range Of Cells Using OpenpyxlCalling an external command in PythonWhat are metaclasses in Python?What is the difference between @staticmethod and @classmethod?Finding the index of an item given a list containing it in PythonDifference between append vs. extend list methods in PythonHow can I safely create a nested directory in Python?Does Python have a ternary conditional operator?Understanding slice notationUnderstanding Python super() with __init__() methodsDoes Python have a string 'contains' substring method?

          How can I change the color of pagination dots of UIPageControl?How to change UIPageControl dotsIs there a way to change page indicator dots colorCustomize dot with image of UIPageControl at index 0 of UIPageControlNo visible @interface for 'NSObject<PageControlDelegate>' declares the selector 'pageControlPageDidChange:'How to change the color of pagination dots in UIPageControl with a different color per pagepagecontrol indicator custom image instead of DefaultChanging the colour of UIPageControl dots in MonoTouchpagecontrol selectable page visibility color?Alternative way to load ViewControllers on a UIPageControlHow to set only layer.border-color for UIpage control dots in swiftHow can I develop for iPhone using a Windows development machine?How to change the name of an iOS app?UITableView - change section header coloruipagecontrol indicator(dot)issueCustom UIPageControl dots color not changingchange the interspace between UIPageControl dotsHow to change Status Bar text color in iOSHow can I change image tintColor in iOS and WatchKitUIPageControl dots with larger space in between each dotsHow to change the color of pagination dots in UIPageControl with a different color per page