【发布时间】: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