【问题标题】:How to catch java.net.ConnectException: Connection refused on akka steram?如何捕获 java.net.ConnectException:akka 流上的连接被拒绝?
【发布时间】:2019-09-10 23:26:28
【问题描述】:

我有一个如下所示的 kafka 消费者:

import akka.actor.ActorSystem
import akka.kafka.scaladsl.Consumer
import akka.kafka.{ConsumerSettings, Subscriptions}
import akka.stream.ActorMaterializer
import akka.stream.scaladsl.Sink
import org.apache.kafka.clients.consumer.ConsumerConfig
import org.apache.kafka.common.serialization.StringDeserializer

import scala.util.{Failure, Success}

object App {
  def main(args: Array[String]): Unit = {


    implicit val system = ActorSystem("SAP-SENDER")
    implicit val executor = system.dispatcher
    implicit val materilizer = ActorMaterializer()

    val config = system.settings.config.getConfig("akka.kafka.consumer")

    val consumerSettings: ConsumerSettings[String, String] =
      ConsumerSettings(config, new StringDeserializer, new StringDeserializer)
        .withBootstrapServers("localhost:9003")
        .withGroupId("SAPSENDER")
        .withProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest")

    Consumer
      .plainSource(
        consumerSettings,
        Subscriptions.topics("TEST-TOPIC")
      )
      .runWith(Sink.foreach(println))
      .onComplete{
        case Success(_) => println("Goood")
        case Failure(ex) =>
          println(s"I am failed ==============> ${ex.getMessage}")
          system.terminate()
      }

  }
} 

kafka 服务器未激活,我只想终止消费者。它总是尝试连接并显示以下消息:

19:03:47.342 [SAP-SENDER-akka.kafka.default-dispatcher-15] DEBUG org.apache.kafka.clients.consumer.KafkaConsumer - [Consumer clientId=consumer-1, groupId=SAPSENDER] Pausing partitions []
19:03:47.342 [SAP-SENDER-akka.kafka.default-dispatcher-15] DEBUG org.apache.kafka.clients.consumer.internals.AbstractCoordinator - [Consumer clientId=consumer-1, groupId=SAPSENDER] No broker available to send FindCoordinator request
19:03:47.342 [SAP-SENDER-akka.kafka.default-dispatcher-15] DEBUG org.apache.kafka.clients.NetworkClient - [Consumer clientId=consumer-1, groupId=SAPSENDER] Give up sending metadata request since no node is available
19:03:47.342 [SAP-SENDER-akka.kafka.default-dispatcher-15] DEBUG org.apache.kafka.clients.consumer.internals.AbstractCoordinator - [Consumer clientId=consumer-1, groupId=SAPSENDER] Coordinator discovery failed, refreshing metadata
19:03:47.342 [SAP-SENDER-akka.kafka.default-dispatcher-15] DEBUG org.apache.kafka.clients.NetworkClient - [Consumer clientId=consumer-1, groupId=SAPSENDER] Give up sending metadata request since no node is available
19:03:47.412 [SAP-SENDER-akka.kafka.default-dispatcher-17] DEBUG org.apache.kafka.clients.consumer.KafkaConsumer - [Consumer clientId=consumer-1, groupId=SAPSENDER] Pausing partitions []
19:03:47.412 [SAP-SENDER-akka.kafka.default-dispatcher-17] DEBUG org.apache.kafka.clients.consumer.internals.AbstractCoordinator - [Consumer clientId=consumer-1, groupId=SAPSENDER] No broker available to send FindCoordinator request
19:03:47.412 [SAP-SENDER-akka.kafka.default-dispatcher-17] DEBUG org.apache.kafka.clients.NetworkClient - [Consumer clientId=consumer-1, groupId=SAPSENDER] Give up sending metadata request since no node is available
19:03:47.412 [SAP-SENDER-akka.kafka.default-dispatcher-17] DEBUG org.apache.kafka.clients.consumer.internals.AbstractCoordinator - [Consumer clientId=consumer-1, groupId=SAPSENDER] Coordinator discovery failed, refreshing metadata
19:03:47.412 [SAP-SENDER-akka.kafka.default-dispatcher-17] DEBUG org.apache.kafka.clients.NetworkClient - [Consumer clientId=consumer-1, groupId=SAPSENDER] Give up sending metadata request since no node is available
19:03:47.478 [SAP-SENDER-akka.kafka.default-dispatcher-20] DEBUG org.apache.kafka.clients.consumer.KafkaConsumer - [Consumer clientId=consumer-1, groupId=SAPSENDER] Pausing partitions []   

它还说:

java.net.ConnectException: Connection refused
    at sun.nio.ch.SocketChannelImpl.checkConnect(Native Method)
    at sun.nio.ch.SocketChannelImpl.finishConnect(SocketChannelImpl.java:717)
    at org.apache.kafka.common.network.PlaintextTransportLayer.finishConnect(PlaintextTransportLayer.java:50)
    at org.apache.kafka.common.network.KafkaChannel.finishConnect(KafkaChannel.java:173)
    at org.apache.kafka.common.network.Selector.pollSelectionKeys(Selector.java:515)
    at org.apache.kafka.common.network.Selector.poll(Selector.java:467)
    at org.apache.kafka.clients.NetworkClient.poll(NetworkClient.java:535)
    at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:265)
    at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:236)
    at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:215)
    at org.apache.kafka.clients.consumer.internals.AbstractCoordinator.ensureCoordinatorReady(AbstractCoordinator.java:231)
    at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.poll(ConsumerCoordinator.java:316)
    at org.apache.kafka.clients.consumer.KafkaConsumer.updateAssignmentMetadataIfNeeded(KafkaConsumer.java:1214)
    at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1179)
    at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1164)
    at akka.kafka.internal.KafkaConsumerActor.poll(KafkaConsumerActor.scala:380)
    at akka.kafka.internal.KafkaConsumerActor.akka$kafka$internal$KafkaConsumerActor$$receivePoll(KafkaConsumerActor.scala:360)
    at akka.kafka.internal.KafkaConsumerActor$$anonfun$receive$1.applyOrElse(KafkaConsumerActor.scala:221)
    at akka.actor.Actor.aroundReceive(Actor.scala:539)
    at akka.actor.Actor.aroundReceive$(Actor.scala:537)
    at akka.kafka.internal.KafkaConsumerActor.akka$actor$Timers$$super$aroundReceive(KafkaConsumerActor.scala:142)
    at akka.actor.Timers.aroundReceive(Timers.scala:51)
    at akka.actor.Timers.aroundReceive$(Timers.scala:40)
    at akka.kafka.internal.KafkaConsumerActor.aroundReceive(KafkaConsumerActor.scala:142)
    at akka.actor.ActorCell.receiveMessage(ActorCell.scala:610)
    at akka.actor.ActorCell.invoke(ActorCell.scala:579)
    at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:268)
    at akka.dispatch.Mailbox.run(Mailbox.scala:229)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)  

如何赶上流中的ConnectException 并阻止消费者尝试连接kafka。

代码托管在这里https://gitlab.com/akka-samples/kafkaconsumer

【问题讨论】:

    标签: scala apache-kafka akka-stream alpakka


    【解决方案1】:

    看着这个PR 和升级到 kafka 客户端 2.0 的工作,我猜想很多重试责任已经委派给了 kafka 客户端。例如,我尝试传递这些属性

    val consumerSettings: ConsumerSettings[String, String] =
      ConsumerSettings(config, new StringDeserializer, new StringDeserializer)
        .withProperties(
          "reconnect.backoff.ms" -> "10000",
          "reconnect.backoff.max.ms" -> "20000"
        )
        .withBootstrapServers("localhost:9099")
        .withGroupId("SAPSENDER")
        .withProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest")
    

    并且异常在 10 秒后第二次出现。我找到了这些属性here

    鉴于此,我认为 kafka 客户端的新适配可能缺少一个功能,因为 KafkaConsumerActor 不会将异常暴露给流,我使用您的 repo 尝试了各种组合,但我仍然得到连续的调试流消息。

    我希望这能给正确的方向一些提示,如果你解决了请告诉我们。

    【讨论】:

      【解决方案2】:

      使用 Kafka Client 2.0+ Alpakka Kafka 无法注意到给定地址上没有可用的 Kafka 代理。

      https://github.com/akka/alpakka-kafka/issues/674

      【讨论】:

        【解决方案3】:

        您应该监控您的流并在出现错误时重新启动它。例如,您可以在 Actor 内部运行您的流,并通过 Actor 监督处理错误连接。

        连接错误可能会持续几秒钟(可能网络不堪重负),因此您应该使用退避策略来避免重试风暴。

        Akka 流已经为您提供了一种使用RestartSource 为流执行此操作的简单方法。见Error Handling

        val control = new AtomicReference[Consumer.Control](Consumer.NoopControl)
        
        val result = RestartSource
          .onFailuresWithBackoff(
            minBackoff = 3.seconds,
            maxBackoff = 30.seconds,
            randomFactor = 0.2
          ) { () =>
            Consumer
              .plainSource(consumerSettings, Subscriptions.topics(topic))
              // this is a hack to get access to the Consumer.Control
              // instances of the latest Kafka Consumer source
              .mapMaterializedValue(c => control.set(c))
              .via(businessFlow)
          }
          .runWith(Sink.seq)
        
        control.get().shutdown()
        

        此解决方案仅在您启动流并且代理关闭时才有效,因为当您尝试创建它时消费者会抛出异常。 但是,如果您成功创建了消费者,然后整个 kafka 集群崩溃,内部 KafkaConsumer 将使用提到的reconnect.backoff.msreconnect.backoff.max.ms 配置重新连接,并且您的流不会失败。

        如果您想限制退休人数,您应该执行以下操作

        val result: Future[Done] = RestartSource
          .onFailuresWithBackoff(
            minBackoff = 3.seconds,
            maxBackoff = 30.seconds,
            randomFactor = 0.2
          ) { () => // your consumer  
          }.
          .take(3) // retries limit
          .runWith(Sink.ignore)
        
        result.onComplete {
          case _ => println("Max retries reached")
        }
        

        【讨论】:

        • 上面显示的代码,我应该在哪里调用control.get().shutdown()。最大后。重试 4 次,我想打电话给control.get().shutdown() 但我应该在哪里打电话呢?
        • 当重试次数用完后,是否有可能将回调传递给RestartSource
        • 我创建了gist.github.com/bifunctor/55440b7515c479f15ddbf8ebd1656e85,但它没有正确停止。 kafka 服务器未激活。
        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2017-08-08
        • 1970-01-01
        • 1970-01-01
        • 2021-05-15
        • 2019-08-20
        • 2017-04-24
        • 1970-01-01
        相关资源
        最近更新 更多