This is an automated email from the ASF dual-hosted git repository. mdedetrich pushed a commit to branch main in repository https://gitbox.apache.org/repos/asf/incubator-pekko-connectors.git
commit f8046eaaf8d43607f3ddd16dc92d85d881bd7791 Author: João Ferreira <[email protected]> AuthorDate: Mon Aug 21 11:25:05 2023 +0100 remove unnecessary Source.fromMaterializer --- .../connectors/kinesis/scaladsl/KinesisSchedulerSource.scala | 8 +------- 1 file changed, 1 insertion(+), 7 deletions(-) diff --git a/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/scaladsl/KinesisSchedulerSource.scala b/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/scaladsl/KinesisSchedulerSource.scala index ebb02527a..c6d4691c4 100644 --- a/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/scaladsl/KinesisSchedulerSource.scala +++ b/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/scaladsl/KinesisSchedulerSource.scala @@ -15,7 +15,6 @@ package org.apache.pekko.stream.connectors.kinesis.scaladsl import org.apache.pekko import pekko.NotUsed -import pekko.dispatch.ExecutionContexts import pekko.stream._ import pekko.stream.connectors.kinesis.impl.KinesisSchedulerSourceStage import pekko.stream.connectors.kinesis.{ @@ -45,12 +44,7 @@ object KinesisSchedulerSource { def apply( schedulerBuilder: ShardRecordProcessorFactory => Scheduler, settings: KinesisSchedulerSourceSettings): Source[CommittableRecord, Future[Scheduler]] = - Source - .fromMaterializer { (mat, _) => - Source - .fromGraph(new KinesisSchedulerSourceStage(settings, schedulerBuilder)) - } - .mapMaterializedValue(_.flatMap(identity)(ExecutionContexts.parasitic)) + Source.fromGraph(new KinesisSchedulerSourceStage(settings, schedulerBuilder)) def sharded( schedulerBuilder: ShardRecordProcessorFactory => Scheduler, --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
