【发布时间】:2020-01-03 18:20:58
【问题描述】:
我有一个简单的可提交源,用于包装在 RestartSource 中的 Kafka 流。它在愉快的路径中运行良好,但如果我故意严重连接到 Kafka 集群,它会从底层 kafka 客户端抛出连接异常并报告 Kafka Consumer Shut Down。我的期望是在约 150 秒后重新启动流,但事实并非如此。从下面我对 RestartSource 的理解/使用是否不正确:
val atomicControl = new AtomicReference[Consumer.Control](NoopControl)
val restartablekafkaSourceWithFlow = {
RestartSource.withBackoff(30.seconds, 120.seconds, 0.2) {
() => {
Consumer.committableSource(consumerSettings.withClientId("clientId"), Subscriptions.topics(Set("someTopic")))
.mapMaterializedValue(c => atomicControl.set(c))
.via(someFlow)
.via(httpFlow)
}
}
}
val committerSink: Sink[(Any, ConsumerMessage.CommittableOffset), Future[Done]] = Committer.sinkWithOffsetContext(CommitterSettings(actorSystem))
val runnableGraph = restartablekafkaSourceWithFlow.toMat(committerSink)(Keep.both)
val control = runnableGraph.mapMaterializedValue(x => Consumer.DrainingControl.apply(atomicControl.get, x._2)).run()
【问题讨论】:
-
所以看来我可以通过 runnableGraph.withAttributes(ActorAttributes.supervisionStrategy(decider)) 对图添加监督,如果底层 kafka 客户端出现连接异常,则重新启动整个图。我不确定为什么上述来源没有失败?
-
您应该提供有关异常的更多详细信息。如果是网络问题,提交也可能失败(毕竟提交者也应该与 kafka 通信),但您的代码仅涵盖源代码部分。如果它真的与接收器连接,您可以尝试使用
Committer.flowWithOffsetContext或将提交接收器包装在RestartSink中。
标签: apache-kafka akka akka-stream alpakka