【问题标题】:Create Source from Lagom/Akka Kafka Topic Subscriber for Websocket从 Lagom/Akka Kafka 主题订阅者为 Websocket 创建源
【发布时间】:2019-02-17 15:33:07
【问题描述】:

我希望我的 Lagom 订阅服务订阅 Kafka 主题并将消息流式传输到 websocket。我使用此文档 (https://www.lagomframework.com/documentation/1.4.x/scala/MessageBrokerApi.html#Subscribe-to-a-topic) 作为指导定义了如下服务:

    // service call
    def stream(): ServiceCall[Source[String, NotUsed], Source[String, NotUsed]]

    // service implementation
    override def stream() = ServiceCall { req =>
      req.runForeach(str => log.info(s"client: %str"))
      kafkaTopic().subscribe.atLeastOnce(Flow.fromFunction(
        // add message to a Source and return Done
      ))
      Future.successful(//some Source[String, NotUsed])

但是,我不太清楚如何处理我的 kafka 消息。 Flow.fromFunction 返回 [String, Done, _] 并暗示我需要将这些消息(字符串)添加到在订阅者之外创建的 Source。

所以我的问题是双重的: 1)如何创建一个akka流源以在运行时接收来自kafka主题订阅者的消息? 2) 如何在 Flow 中将 kafka 消息附加到所述源?

【问题讨论】:

  • req.runForeach(str => log.info(s"client: %str")) 将耗尽您的整个输入流。

标签: scala websocket apache-kafka akka-stream lagom


【解决方案1】:

您似乎误解了 Lagom 的服务 API。如果你试图从你的服务调用主体中实现一个流,你的调用没有输入;即,

def stream(): ServiceCall[Source[String, NotUsed], Source[String, NotUsed]]

暗示当客户端提供Source[String, NotUsed]时,服务会以实物回应。您的客户没有直接提供此服务;因此,您的签名应该是

def stream(): ServiceCall[NotUsed, Source[String, NotUsed]]

现在回答你的问题...

这实际上在 scala giter8 模板中不存在,但 java 版本包含他们称之为 autonomous stream 的东西,它大致完成了你想做的事情。

在 Scala 中,这段代码看起来像...

override def autonomousStream(): ServiceCall[
  Source[String, NotUsed], 
  Source[String, NotUsed]
] = ServiceCall { hellos => Future {
    hellos.mapAsync(8, ...)
  }
}

由于您的调用不是映射到 input 流,而是映射到 kafka 主题,因此您需要执行以下操作:

override def stream(): ServiceCall[NotUsed, Source[String, NotUsed]] = ServiceCall { 
  _ => 
    Future {
      kafkaTopic()
        .subscribe
        .atMostOnce
        .mapAsync(...)
    }
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2013-07-03
    • 2017-07-31
    • 2020-10-29
    • 1970-01-01
    • 2020-01-12
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多