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]