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]
