【问题标题】:Kafka message to websocketKafka消息到websocket
【发布时间】:2016-04-01 04:25:56
【问题描述】:

我正在尝试使用 reactive-kafka、akka-http 和 akka-stream 编写一个 Kafka 消费者到 websocket 流。

  val publisherActor = actorSystem.actorOf(CommandPublisher.props)
  val publisher = ActorPublisher[String](publisherActor)
  val commandSource = Source.fromPublisher(publisher) map toMessage
  def toMessage(c: String): Message = TextMessage.Strict(c)

  class CommandPublisher extends ActorPublisher[String] {
    override def receive = {
      case cmd: String =>
        if (isActive && totalDemand > 0)
          onNext(cmd)
    }
  }

  object CommandPublisher {
    def props: Props = Props(new CommandPublisher())
  }

  // This is the route 
  def mainFlow(): Route = {
    path("ws" / "commands" ) {
       handleWebSocketMessages(Flow.fromSinkAndSource(Sink.ignore, commandSource))
    } 
  }

从 kafka 消费者(这里省略),我做了一个publisherActor ! commandString 来动态添加内容到 websocket。

但是,当我启动多个客户端到 websocket 时,我在后端遇到了这个异常:

[ERROR] [03/31/2016 21:17:10.335] [KafkaWs-akka.actor.default-dispatcher-3][akka.actor.ActorSystemImpl(KafkaWs)] WebSocket handler failed with can not subscribe the same subscriber multiple times (see reactive-streams specification, rules 1.10 and 2.12)
java.lang.IllegalStateException: can not subscribe the same subscriber multiple times (see reactive-streams specification, rules 1.10 and 2.12)
  at akka.stream.impl.ReactiveStreamsCompliance$.canNotSubscribeTheSameSubscriberMultipleTimesException(ReactiveStreamsCompliance.scala:35)
  at akka.stream.actor.ActorPublisher$class.aroundReceive(ActorPublisher.scala:295)
  ...

一个流不能用于所有 websocket 客户端吗?还是应该为每个客户端创建流/发布者角色?

在这里,我打算向所有 websocket 客户端发送“当前”/“实时”通知。通知历史无关紧要,新客户需要忽略。

【问题讨论】:

    标签: websocket akka-stream akka-http


    【解决方案1】:

    我很抱歉承受坏消息,但看起来这是akka 对 .您不能随意为所有客户端重用流的实例。由于 Rx 模型,扇出必须是“显式”的。

    我遇到的示例使用特定于路由的Flow

      // The flow from beginning to end to be passed into handleWebsocketMessages
      def websocketDispatchFlow(sender: String): Flow[Message, Message, Unit] =
        Flow[Message]
          // First we convert the TextMessage to a ReceivedMessage
          .collect { case TextMessage.Strict(msg) => ReceivedMessage(sender, msg) }
          // Then we send the message to the dispatch actor which fans it out
          .via(dispatchActorFlow(sender))
          // The message is converted back to a TextMessage for serialization across the socket
          .map { case ReceivedMessage(from, msg) => TextMessage.Strict(s"$from: $msg") }
    
      def route =
        (get & path("chat") & parameter('name)) { name =>
          handleWebsocketMessages(websocketDispatchFlow(sender = name))
        }
    

    下面是关于它的讨论:

    这正是我在 Akka Stream 中不喜欢的,这种明确的 扇出。当我从某个地方收到我想要的数据源时 进程(例如 Observable 或 Source),我只想订阅它 我不想在乎它是冷是热还是 是否已被其他订阅者订阅。这是我的河流类比。 河流不应该在乎谁喝它和喝水的人 不应该关心河流的源头或其他多少 有饮酒者。我的样本,相当于一个 Mathias 提供,确实共享数据源,但它只是引用 计数,你可以有 2 个订阅者,或者你可以有 100 个,不是 事情。在这里我很喜欢,但引用计数没有 如果您不想丢失事件或者如果您想确保 流始终保持开启状态。但是你使用ConnectableObservable 其中有connect(): Cancelable,非常适合说... Play 的 LifeCycle 插件。并且在此基础上您可以使用 BehaviorSubject 或 ReplaySubject 如果你想重复上一个 新订户的价值。之后事情就正常了,没有手册 需要绘制该连接图。 ... ...(来自https://bionicspirit.com/blog/2015/09/06/monifu-vs-akka-streams.html) ... 对于接受 Observable 并返回 Observable 的函数,我们 确实有电梯,这是最接近具有 名称,可用于 Subject 或其他 由于 LiftOperators1(和 2)的可观察类型,这就是 可以在不丢失类型的情况下转换 Observables - 这是对 RxJava 对 lift 所做的 OOP 的改进。

    但是,这样的功能并不等同于Processor / Subject。这 区别在于Subject 同时是消费者和 生产者。这意味着订阅者无法完全控制 当数据源启动并且数据源本质上是 hot(意味着多个订阅者共享同一个数据源)。在 Rx 中,如果你对 cold observables 建模(意思是 为每个人启动一个新数据源的 observables 订户)。另一方面,在 Rx 中(一般而言),拥有 只能订阅一次的数据源,仅此而已。这 Monifu 中这条规则的唯一例外是由 GroupBy 运算符,但这就像确认 规则。

    这意味着什么,尤其是加上另一个限制 Monifu 和 Reactive Streams 协议的合同(您应 不与同一个消费者多次订阅),是 SubjectProcessor 实例不可重用。为了这样 一个可重用的实例,Rx 模型需要一个工厂 Processor。此外,这意味着每当您想使用 Subject / Processor,您的数据源必须自动 (可在多个订阅者之间共享)。

    【讨论】:

    • 有趣的是,我采用了来自github.com/J-Technologies/… 的代码,这出乎意料地有效。它使用 Source.actorPublisher 而不是 Source.fromPublisher - 这是唯一的区别。无法理解为什么会这样,而我的不行
    • 我很确定我引用的帖子是准确的......所以这意味着,你必须像使用普通 Actor 一样使用它,而不是与 ReactiveStreams 兼容
    • 来自ActorPublisher 它说:/** * 创建一个由 [[ActorPublisher]] 演员支持的 [[org.reactivestreams.Publisher]]。它可以 * 附加到 [[org.reactivestreams.Subscriber]] 或用作 * [[akka.stream.scaladsl.Flow]] 的输入源。 */
    猜你喜欢
    • 2023-03-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-06-06
    • 2016-05-15
    • 2020-06-04
    相关资源
    最近更新 更多