【发布时间】:2020-03-05 16:08:40
【问题描述】:
我目前正在运行类似于以下内容的 Akka 流设置:
┌───────────────┐
┌─────────────┐ │┌─────────────┐│
│REST endpoint│──▶│Queue source ││
└─────────────┘ │└──────╷──────┘│
│┌──────▼──────┐│
││ Flow[T] ││
│└──────╷──────┘│
│┌──────▼──────┐│ ┌─────────────┐
││ KafkaSink │├─▶│ Kafka topic │
│└─────────────┘│ └─────────────┘
└───────────────┘
虽然这项工作运行良好,但我想对生产系统有一些了解,即是否存在错误以及错误类型。例如,我将KafkaSink 包装成RestartSink.withBackoff 并将以下属性应用于包装的接收器:
private val decider: Supervision.Decider = {
case x =>
log.error("KafkaSink encountered an error and will stop", x)
Supervision.Stop
}
Flow[...]
.log("KafkaSink")
.to(Producer.plainSink(...))
.withAttributes(ActorAttributes.supervisionStrategy(decider))
.addAttributes(
ActorAttributes.logLevels(
onElement = Logging.DebugLevel,
onFinish = Logging.WarningLevel,
onFailure = Logging.ErrorLevel
)
)
这确实为我提供了一些见解,例如我将收到一条日志消息,说明下游已关闭,以及通过我添加的 supervisionStrategy 发生的异常。
然而,这个解决方案感觉有点像一种变通方法(例如,将异常记录在监督策略中),并且它也没有提供任何对RestartWithBackoffSink 行为的洞察。当然,我可以为该类启用DEBUG 级别的日志记录,但同样,这感觉像是在生产环境中做的一种解决方法。
长话短说:
- 我试图深入了解 Akka 流中发生的错误的方式是否有任何明显的缺点
- 是否有更好/更惯用的方式在生产环境中向 Akka 流添加日志记录
【问题讨论】:
标签: scala logging akka akka-stream