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 26674de73c116bf98532d3b371efeda4d6678f16
Author: João Ferreira <[email protected]>
AuthorDate: Fri Aug 18 15:19:26 2023 +0100

    kinesis: use stage materializer with IODispatcher instead of injected 
execution context
---
 .../stream/connectors/kinesis/impl/KinesisSchedulerSourceStage.scala  | 4 +++-
 .../stream/connectors/kinesis/scaladsl/KinesisSchedulerSource.scala   | 1 -
 2 files changed, 3 insertions(+), 2 deletions(-)

diff --git 
a/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/KinesisSchedulerSourceStage.scala
 
b/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/KinesisSchedulerSourceStage.scala
index f30e11302..0518f3d0c 100644
--- 
a/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/KinesisSchedulerSourceStage.scala
+++ 
b/kinesis/src/main/scala/org/apache/pekko/stream/connectors/kinesis/impl/KinesisSchedulerSourceStage.scala
@@ -49,7 +49,7 @@ private[kinesis] object KinesisSchedulerSourceStage {
 @InternalApi
 private[kinesis] class KinesisSchedulerSourceStage(
     settings: KinesisSchedulerSourceSettings,
-    schedulerBuilder: ShardRecordProcessorFactory => Scheduler)(implicit ec: 
ExecutionContext)
+    schedulerBuilder: ShardRecordProcessorFactory => Scheduler)
     extends GraphStageWithMaterializedValue[SourceShape[CommittableRecord], 
Future[Scheduler]] {
 
   private val out = Outlet[CommittableRecord]("Records")
@@ -76,6 +76,8 @@ private[kinesis] class KinesisSchedulerSourceStage(
     private[this] val buffer = mutable.Queue.empty[CommittableRecord]
     private[this] var schedulerOpt: Option[Scheduler] = None
 
+    implicit def ec: ExecutionContext = materializer.executionContext
+
     override def preStart(): Unit = {
       val scheduler = schedulerBuilder(new ShardRecordProcessorFactory {
         override def shardRecordProcessor(): ShardRecordProcessor =
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 53d0fc4e1..ebb02527a 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
@@ -47,7 +47,6 @@ object KinesisSchedulerSource {
       settings: KinesisSchedulerSourceSettings): Source[CommittableRecord, 
Future[Scheduler]] =
     Source
       .fromMaterializer { (mat, _) =>
-        import mat.executionContext
         Source
           .fromGraph(new KinesisSchedulerSourceStage(settings, 
schedulerBuilder))
       }


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to