This is an automated email from the ASF dual-hosted git repository.
pjfanning pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko-management.git
The following commit(s) were added to refs/heads/main by this push:
new 3286db00 improve discovery-aws-api (#907)
3286db00 is described below
commit 3286db00434d0c6aacd67d8305e5dc4bf2575109
Author: PJ Fanning <[email protected]>
AuthorDate: Sun Aug 9 09:40:53 2026 +0100
improve discovery-aws-api (#907)
* improve discovery-aws-api
* use pekko.actor.default-blocking-io-dispatcher for blocking tasks
* rework shutdown based on review
---
.../awsapi/ec2/Ec2TagBasedServiceDiscovery.scala | 23 ++++++++++++++++++----
.../discovery/awsapi/ecs/EcsServiceDiscovery.scala | 21 +++++++++++++++++---
2 files changed, 37 insertions(+), 7 deletions(-)
diff --git
a/discovery-aws-api/src/main/scala/org/apache/pekko/discovery/awsapi/ec2/Ec2TagBasedServiceDiscovery.scala
b/discovery-aws-api/src/main/scala/org/apache/pekko/discovery/awsapi/ec2/Ec2TagBasedServiceDiscovery.scala
index f703a971..c5064d8c 100644
---
a/discovery-aws-api/src/main/scala/org/apache/pekko/discovery/awsapi/ec2/Ec2TagBasedServiceDiscovery.scala
+++
b/discovery-aws-api/src/main/scala/org/apache/pekko/discovery/awsapi/ec2/Ec2TagBasedServiceDiscovery.scala
@@ -19,7 +19,7 @@ import com.amazonaws.retry.PredefinedRetryPolicies
import com.amazonaws.services.ec2.model.{ DescribeInstancesRequest, Filter,
Reservation }
import com.amazonaws.services.ec2.{ AmazonEC2, AmazonEC2ClientBuilder }
import org.apache.pekko
-import pekko.actor.ExtendedActorSystem
+import pekko.actor.{ CoordinatedShutdown, ExtendedActorSystem }
import pekko.annotation.InternalApi
import pekko.discovery.ServiceDiscovery.{ Resolved, ResolvedTarget }
import
pekko.discovery.awsapi.ec2.Ec2TagBasedServiceDiscovery.parseFiltersString
@@ -100,7 +100,10 @@ final class Ec2TagBasedServiceDiscovery(system:
ExtendedActorSystem) extends Ser
}
}
- private val ec2Client: AmazonEC2 = {
+ @volatile private var ec2ClientUsed = false
+
+ private lazy val ec2Client: AmazonEC2 = {
+ ec2ClientUsed = true
val clientConfiguration = clientConfigFqcn match {
case Some(fqcn) =>
getCustomClientConfigurationInstance(fqcn) match {
@@ -130,6 +133,17 @@ final class Ec2TagBasedServiceDiscovery(system:
ExtendedActorSystem) extends Ser
builder.build()
}
+ CoordinatedShutdown(system).addTask(CoordinatedShutdown.PhaseServiceUnbind,
"ec2-client-close") { () =>
+ if (ec2ClientUsed) {
+ Future {
+ ec2Client.shutdown()
+ pekko.Done
+ }(ec)
+ } else {
+ Future.successful(pekko.Done)
+ }
+ }
+
@tailrec
private def getInstances(
client: AmazonEC2,
@@ -144,9 +158,10 @@ final class Ec2TagBasedServiceDiscovery(system:
ExtendedActorSystem) extends Ser
val describeInstancesResult =
client.describeInstances(describeInstancesRequest)
val ips: List[String] =
- describeInstancesResult.getReservations.asScala.toList
- .flatMap((r: Reservation) => r.getInstances.asScala.toList)
+ describeInstancesResult.getReservations.asScala
+ .flatMap((r: Reservation) => r.getInstances.asScala)
.map(instance => instance.getPrivateIpAddress)
+ .toList
val accumulatedIps = accumulator ++ ips
diff --git
a/discovery-aws-api/src/main/scala/org/apache/pekko/discovery/awsapi/ecs/EcsServiceDiscovery.scala
b/discovery-aws-api/src/main/scala/org/apache/pekko/discovery/awsapi/ecs/EcsServiceDiscovery.scala
index e7c588f1..ad84469f 100644
---
a/discovery-aws-api/src/main/scala/org/apache/pekko/discovery/awsapi/ecs/EcsServiceDiscovery.scala
+++
b/discovery-aws-api/src/main/scala/org/apache/pekko/discovery/awsapi/ecs/EcsServiceDiscovery.scala
@@ -17,7 +17,7 @@ import java.net.{ InetAddress, NetworkInterface }
import java.util.concurrent.TimeoutException
import org.apache.pekko
-import pekko.actor.ActorSystem
+import pekko.actor.{ ActorSystem, CoordinatedShutdown }
import pekko.discovery.{ Lookup, ServiceDiscovery }
import pekko.discovery.ServiceDiscovery.{ Resolved, ResolvedTarget }
import pekko.discovery.awsapi.ecs.EcsServiceDiscovery.resolveTasks
@@ -48,7 +48,10 @@ final class EcsServiceDiscovery(system: ActorSystem) extends
ServiceDiscovery {
private val config =
system.settings.config.getConfig("pekko.discovery.aws-api-ecs")
private val cluster = config.getString("cluster")
+ @volatile private var ecsClientUsed = false
+
private lazy val ecsClient = {
+ ecsClientUsed = true
// we have our own retry/backoff mechanism, so we don't need EC2Client's
in addition
val clientConfiguration = new ClientConfiguration()
clientConfiguration.setRetryPolicy(PredefinedRetryPolicies.NO_RETRY_POLICY)
@@ -63,7 +66,19 @@ final class EcsServiceDiscovery(system: ActorSystem) extends
ServiceDiscovery {
builder.build()
}
- private implicit val ec: ExecutionContext = system.dispatcher
+ private implicit val ec: ExecutionContext =
+ system.dispatchers.lookup("pekko.actor.default-blocking-io-dispatcher")
+
+ CoordinatedShutdown(system).addTask(CoordinatedShutdown.PhaseServiceUnbind,
"ecs-client-close") { () =>
+ if (ecsClientUsed) {
+ Future {
+ ecsClient.shutdown()
+ pekko.Done
+ }(ec)
+ } else {
+ Future.successful(pekko.Done)
+ }
+ }
override def lookup(query: Lookup, resolveTimeout: FiniteDuration):
Future[Resolved] =
Future.firstCompletedOf(
@@ -144,7 +159,7 @@ object EcsServiceDiscovery {
private def describeTasks(ecsClient: AmazonECS, cluster: String, taskArns:
Seq[String]): Seq[Task] =
for {
// Each DescribeTasksRequest can contain at most 100 task ARNs.
- group <- taskArns.grouped(100).toList
+ group <- taskArns.grouped(100).toSeq
tasks = ecsClient.describeTasks(
new
DescribeTasksRequest().withCluster(cluster).withTasks(group.asJava))
task <- tasks.getTasks.asScala
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]