【问题标题】:Kafka topic to websocketKafka 主题到 websocket
【发布时间】:2017-07-27 04:14:58
【问题描述】:

我正在尝试实现一个设置,其中我有多个 Web 浏览器打开到我的 akka-http 服务器的 websocket 连接,以便读取发布到 kafka 主题的所有消息。

所以消息流应该是这样的

kafka topic -> akka-http -> websocket connection 1 
                         -> websocket connection 2
                         -> websocket connection 3 

现在我已经为 websocket 创建了一个路径:

val route: Route = 
 path("ws") {
   handleWebSocketMessages(notificationWs)
 }

然后我为我的 kafka 主题创建了一个消费者:

val consumerSettings = ConsumerSettings(system,
  new ByteArrayDeserializer, new StringDeserializer)
    .withBootstrapServers("localhost:9092")
    .withGroupId("group1")
    .withProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest")
val source = Consumer
  .plainSource(consumerSettings, Subscriptions.topics("topic1"))

最后我想把这个源连接到handleWebSocketMessages中的websocket

def handleWebSocketMessages: Flow[Message, Message, Any] =
  Flow[Message].mapConcat {
    case tm: TextMessage =>
      TextMessage(source)::Nil
    case bm: BinaryMessage =>
      // ignore binary messages but drain content to avoid the stream being clogged
      bm.dataStream.runWith(Sink.ignore)
      Nil
  }

这是我尝试在 TextMessage 中使用 source 时遇到的错误:

错误:(77, 9) 重载方法值适用于替代方案: (textStream: akka.stream.scaladsl.Source[String,Any])akka.http.scaladsl.model.ws.TextMessage (文本:字符串)akka.http.scaladsl.model.ws.TextMessage.Strict 不能应用于 (akka.stream.scaladsl.Source[org.apache.kafka.clients.consumer.ConsumerRecord[Array[Byte],String],akka.kafka.scaladsl.Consumer.Control]) TextMessage(source)::Nil

我认为我在此过程中犯了很多错误,但我想说最阻碍的部分是handleWebSocketMessages

【问题讨论】:

    标签: scala apache-kafka akka-stream akka-http


    【解决方案1】:

    首先,要了解源的类型是:Source[ConsumerRecord[K, V], Control]。 因此,这不是您可以作为 TextMessage 的参数传递的东西。

    现在,让我们从 websocket 的角度来看:

    • 为 Kafka 源中的每条消息构建传出消息。该消息将是来自 Kafka 消息的字符串转换的 TextMessage。
    • 对于每条传入的消息,只需 println() 它

    因此,Flow 可以被视为两个组件:SourceSink

    val incomingMessages: Sink[Message, NotUsed] =
      Sink.foreach(println(_))
    
    val outgoingMessages: Source[Message, NotUsed] =
      source
        .map { consumerRecord => TextMessage(consumerRecord.record.value) }
    
    val handleWebSocketMessages: Flow[Message, Message, Any]  
      = Flow.fromSinkAndSource(incomingMessages, outgoingMessages)
    

    希望对您有所帮助。

    【讨论】:

    • 非常感谢!在我使用 Consumer.committableSource 而不是 Consumer.plainSourceconsumerRecord.record.value() 而不是 consumerRecord.getkey.toString 之后,您的答案有效。
    猜你喜欢
    • 1970-01-01
    • 2023-03-16
    • 2021-11-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-02-17
    • 2015-09-03
    相关资源
    最近更新 更多