This is an automated email from the ASF dual-hosted git repository.

raboof pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko-connectors-kafka.git


The following commit(s) were added to refs/heads/main by this push:
     new 28947a28 propagate KafkaConnectionFailed to owner actor (#626)
28947a28 is described below

commit 28947a28284c5cedcba6685ab11786e22023349d
Author: PJ Fanning <[email protected]>
AuthorDate: Mon Aug 10 08:38:46 2026 +0100

    propagate KafkaConnectionFailed to owner actor (#626)
    
    * propagate KafkaConnectionFailed to owner actor
    
    * Update KafkaConsumerActor.scala
---
 .../pekko/kafka/internal/KafkaConsumerActor.scala  |  4 ++-
 .../apache/pekko/kafka/internal/ConsumerSpec.scala | 36 +++++++++++++++++++++-
 2 files changed, 38 insertions(+), 2 deletions(-)

diff --git 
a/core/src/main/scala/org/apache/pekko/kafka/internal/KafkaConsumerActor.scala 
b/core/src/main/scala/org/apache/pekko/kafka/internal/KafkaConsumerActor.scala
index 5f999c16..6faa41ae 100644
--- 
a/core/src/main/scala/org/apache/pekko/kafka/internal/KafkaConsumerActor.scala
+++ 
b/core/src/main/scala/org/apache/pekko/kafka/internal/KafkaConsumerActor.scala
@@ -659,10 +659,12 @@ import scala.util.control.NonFatal
 
   private def processErrors(exception: Throwable): Unit = {
     val sendTo = (stageActorsMap.values ++ owner).toSet
-    log.debug(s"sending failure {} to {}", exception.getClass, 
sendTo.mkString(","))
+    if (log.isDebugEnabled)
+      log.debug("sending failure {} to {}", exception.getClass, 
sendTo.mkString(","))
     stageActorsMap.values.foreach { stageActorRef =>
       sendFailure(exception, stageActorRef)
     }
+    owner.foreach(_ ! Failure(exception))
   }
 
   private def handleMetadataRequest(req: Metadata.Request): Metadata.Response 
= req match {
diff --git 
a/tests/src/test/scala/org/apache/pekko/kafka/internal/ConsumerSpec.scala 
b/tests/src/test/scala/org/apache/pekko/kafka/internal/ConsumerSpec.scala
index b4c762c2..12b32977 100644
--- a/tests/src/test/scala/org/apache/pekko/kafka/internal/ConsumerSpec.scala
+++ b/tests/src/test/scala/org/apache/pekko/kafka/internal/ConsumerSpec.scala
@@ -21,13 +21,17 @@ import pekko.kafka.ConsumerMessage._
 import pekko.kafka.scaladsl.Consumer
 import pekko.kafka.scaladsl.Consumer.Control
 import pekko.kafka.tests.scaladsl.LogCapturing
-import pekko.kafka.{ CommitTimeoutException, ConsumerSettings, Repeated, 
Subscriptions }
+import pekko.actor.Status.Failure
+import pekko.kafka.{ CommitTimeoutException, ConsumerSettings, 
KafkaConnectionFailed, Repeated, Subscriptions }
+import pekko.kafka.{ KafkaConsumerActor => PublicKafkaConsumerActor }
+import pekko.testkit.TestProbe
 import pekko.stream.scaladsl._
 import pekko.stream.testkit.scaladsl.StreamTestKit.assertAllStagesStopped
 import pekko.stream.testkit.scaladsl.TestSink
 import pekko.testkit.TestKit
 import com.typesafe.config.ConfigFactory
 import org.apache.kafka.clients.consumer._
+import org.apache.kafka.common.TopicPartition
 import org.apache.kafka.common.serialization.StringDeserializer
 import org.mockito.Mockito._
 import org.scalatest.BeforeAndAfterAll
@@ -335,4 +339,34 @@ class ConsumerSpec(_system: ActorSystem)
     Await.result(control.isShutdown, remainingOrDefault)
     mock.verifyClosed()
   }
+
+  it should "propagate KafkaConnectionFailed to owner actor" in {
+    val mock = new ConsumerMock[K, V]()
+    val ownerProbe = TestProbe()
+    val settings = ConsumerSettings
+      .create(system, new StringDeserializer, new StringDeserializer)
+      .withGroupId("group1")
+      .withCloseTimeout(ConsumerMock.closeTimeout)
+      .withConsumerFactory(_ => mock.mock)
+    val actor = system.actorOf(PublicKafkaConsumerActor.props(ownerProbe.ref, 
settings))
+
+    // trigger initialization by subscribing
+    actor ! KafkaConsumerActor.Internal.Subscribe(
+      Set("topic"),
+      new org.apache.pekko.kafka.scaladsl.PartitionAssignmentHandler {
+        override def onAssign(assignment: Set[TopicPartition],
+            restrictedConsumer: org.apache.pekko.kafka.RestrictedConsumer): 
Unit = ()
+        override def onRevoke(revokedTps: Set[TopicPartition],
+            restrictedConsumer: org.apache.pekko.kafka.RestrictedConsumer): 
Unit = ()
+        override def onLost(lostTps: Set[TopicPartition],
+            restrictedConsumer: org.apache.pekko.kafka.RestrictedConsumer): 
Unit = ()
+        override def onStop(currentTps: Set[TopicPartition],
+            restrictedConsumer: org.apache.pekko.kafka.RestrictedConsumer): 
Unit = ()
+      })
+
+    val kcf = KafkaConnectionFailed(new 
org.apache.kafka.common.errors.TimeoutException("connection lost"), 3)
+    actor ! kcf
+
+    ownerProbe.expectMsg(Failure(kcf))
+  }
 }


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

Reply via email to