【问题标题】:Why Logging is Not Working for Akka Stream为什么日志记录不适用于 Akka Stream
【发布时间】:2019-05-02 19:06:22
【问题描述】:

我正在使用 Alpakka,下面有玩具示例:

val system = ActorSystem("system")
implicit val materializer: ActorMaterializer = ActorMaterializer.create(system)

implicit val adapter: LoggingAdapter = Logging(system, "customLogger")
implicit val ec: ExecutionContextExecutor = system.dispatcher

val log = Logger(this.getClass, "Foo")

val consumerConfig = system.settings.config.getConfig("akka.kafka.consumer")
val consumerSettings: ConsumerSettings[String, String] =
  ConsumerSettings(consumerConfig, new StringDeserializer, new StringDeserializer)
    .withBootstrapServers("localhost:9092")
    .withGroupId("my-group")

def start() = {
  Consumer.plainSource(consumerSettings, Subscriptions.topics("test"))
    .log("My Consumer: ")
    .withAttributes(
      Attributes.logLevels(
        onElement = Logging.InfoLevel,
        onFinish = Logging.InfoLevel,
        onFailure = Logging.DebugLevel
      )
    )
    .filter(//some predicate)
    .map(// some process)
    .map(out => ByteString(out))
    .runWith(LogRotatorSink(timeFunc))
    .onComplete {
      case Success(_) => log.info("DONE")
      case Failure(e) => log.error("ERROR")
    }
}

此代码正在运行。但我在记录时遇到问题。带有属性的第一部分是日志记录很好。当元素进入时,它会将日志记录到标准输出。但是当 LogRotatorSink 完成并且未来完成时,我想将 DONE 打印到标准输出。这是行不通的。正在生成文件,因此进程正在运行,但没有向标准输出发送“DONE”消息。

请问如何将“DONE”输出到标准输出?

akka {

  # Loggers to register at boot time (akka.event.Logging$DefaultLogger logs
  # to STDOUT)
  loggers = ["akka.event.slf4j.Slf4jLogger"]

  # Log level used by the configured loggers (see "loggers") as soon
  # as they have been started; before that, see "stdout-loglevel"
  # Options: OFF, ERROR, WARNING, INFO, DEBUG
  loglevel = "INFO"

  # Log level for the very basic logger activated during ActorSystem startup.
  # This logger prints the log messages to stdout (System.out).
  # Options: OFF, ERROR, WARNING, INFO, DEBUG
  stdout-loglevel = "INFO"

  # Filter of log events that is used by the LoggingAdapter before
  # publishing log events to the eventStream.
  logging-filter = "akka.event.slf4j.Slf4jLoggingFilter"

}


<configuration>

    <appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
        <encoder>
            <pattern>%highlight(%date{HH:mm:ss.SSS} %-5level %-50.50([%logger{50}])) - %msg%n</pattern>
        </encoder>
    </appender>

    <logger name="org.apache.kafka" level="INFO"/>

    <root level="INFO">
        <appender-ref ref="STDOUT"/>
    </root>

</configuration>

【问题讨论】:

    标签: scala logging akka akka-stream alpakka


    【解决方案1】:

    日志正在工作 - 是你的 Future 没有结束,因为 Kafka Consumer 是一个无限流 - 当它会读取所有内容并到达主题中的最新消息时......它将等待新消息出现- 在许多情况下,例如即使突然关闭此类流的采购将是一场灾难,因此默认无限运行流是明智的选择。

    这个流应该在什么时候结束?明确定义此条件,您将能够使用.take(n).takeUntil(cond).takeWithin(time) 之类的东西在明确定义的条件下关闭它。然后流将关闭,Future 将完成,您的 DONE 将被打印出来。

    【讨论】:

    • 谢谢。我想要的是在处理完每个元素后打印到 STDOUT。因此,当 LogRotatorSink 完成写入文件时,我想打印消息以记录它已在流中完成此元素
    • 谢谢。我知道如何解决问题。我不使用 LogRotatorSink,但我只使用 Sink。然后我执行 .to(Sink.foreach(.....)) - 在 foreach 内部,我将执行日志记录的函数传递到 STDOUT。但我很遗憾我不能使用 LogRotatorSink 并改变他来记录到 STDOUT
    • 实际上,除了foreach(必须是Sink),您还可以.map {a =&gt; log(a); a } 并继续处理流。然后你可以使用LogRotatorSink
    猜你喜欢
    • 2020-02-02
    • 2019-03-01
    • 2019-09-16
    • 1970-01-01
    • 2014-11-23
    • 2018-07-27
    • 1970-01-01
    • 2016-07-02
    • 2019-02-20
    相关资源
    最近更新 更多