【问题标题】:Is it safe to pass an RX Observable to an actor (scala)?将 RX Observable 传递给演员(scala)是否安全?
【发布时间】:2015-09-10 15:51:04
【问题描述】:

我已经为 RX Java 使用 scala 绑定有一段时间了,并且正在考虑将其与 Akka Actors 结合使用。 我想知道在 Akka Actors 之间传递 RX Observables 是否安全/可行。 例如,一个打印最多 20 的偶数平方(每秒)的程序:

/* producer creates an observable and sends it to the worker */
object Producer extends Actor {
  val toTwenty : Observable[Int] = Observable.interval(1 second).take(20)

  def receive = {
    case o : Observable[Int] =>
      o.subscribe( onNext => println )
  }

  worker ! toTwenty
}


/* worker which returns squares of even numbers */
object Worker extends Actor {
  def receive = {
    case o : Observable[Int] => 
       sender ! o filter { _ % 2 == 0 } map { _^2 }
  }
}

(请将此视为伪代码;它不会编译)。注意我是sending Observables 从一个演员到另一个演员。我想了解:

  • Akka 和 RX 会自动同步对 Observable 的访问吗?
  • Observable 不能通过分布式系统发送 - 它是对本地内存中对象的引用。但是,它会在本地工作吗?
  • 假设在这个简单的示例中,工作将安排在subscribe 调用中的Producer 上。我可以拆分工作,以便分别在每个演员身上完成吗?

题外话:我见过一些将 RX 和 Actor 结合起来的项目:

http://jmhofer.johoop.de/?p=507https://github.com/jmhofer/rxjava-akka

但它们的不同之处在于它们不会简单地将Observable 作为参与者之间的消息传递。他们首先调用subscribe() 来获取值,然后将它们发送到演员邮箱,并从中创建一个新的Observable。还是我弄错了?

【问题讨论】:

    标签: scala akka rx-java


    【解决方案1】:

    您的方法不是一个好主意。 Akka 背后的主要思想是将消息发送到 Actor 的邮箱,然后 Actor 按顺序处理它们(在一个线程上)。 这样就不可能有 2 个线程访问一个 Actor 的状态并且不会出现并发问题。

    在您的情况下,您在 Observable 上使用订阅。您的 onNext 回调可能会在另一个线程上执行。因此,突然有可能 2 个线程可以访问您的演员的状态。所以你必须非常小心你在回调中所做的事情。这是您最后一次观察其他实现的原因。这些实现似乎在 onNext 中获取值并将该值作为消息发送。 您不得在此类回调中更改参与者的内部状态。而是向同一个actor发送消息。这样可以再次保证在一个线程上进行顺序处理。

    【讨论】:

    • 感谢您的回复。你介意看看我发布的第二个链接(github),并告诉我自述文件中的第一个代码示例在做什么吗?它以某种方式将 observable 与 actor 同步,因此您可以在 receive 中使用它。这是否允许您在 actor 内部使用 Observable(即使您不能将其作为消息传递)?
    • 嘿,你有一个关于你想用 Observable 做什么的例子吗?第二个 github 链接显示了一个集成,但是这个例子太简单了,因为它只记录日志并且没有任何副作用。
    【解决方案2】:

    我花了一些时间进行实验,发现你可以在 Akka 中使用Observables。事实上,由于Observable 可以被认为是Future 的多变量扩展,因此您可以遵循与组合Actor 和Futures 相同的准则。实际上,官方文档和教科书(例如 Akka Concurrency,Wyatt 2013)都支持/鼓励在 Akka 中使用 Future,但有很多警告。

    首先是积极的:

    • Observables 和 Futures 一样是不可变的,因此理论上它们在消息中传递应该是安全的。
    • Observable 允许您指定执行上下文,非常类似于Future。这是使用Observable.observeOn(scheduler) 完成的。您可以通过将 Akka 调度程序(例如 system.dispatchercontext.dispatcher)传递给 rx.lang.scala.ExecutorScheduler 构造函数,从 Akka 的 exec 上下文创建调度程序。这应该确保它们是同步的。
    • 与上述相关,正在开发中的 rx-scala 增强功能 (https://github.com/Netflix/RxJava/issues/815#issuecomment-38793433),这将允许隐式指定 observable 的调度程序。
    • Futures 非常适合使用 ask 模式的 Akka。类似的模式可以用于 Observables(见本文底部)。这也解决了向远程 observables 发送消息的问题。

    现在需要注意的是:

    • 它们与未来有着相同的问题。例如,请参见页面底部:http://doc.akka.io/docs/akka/2.3.2/general/jmm.html。还有 Wyatt 2013 中关于期货的章节。
    • 正如@mavilein 的回答,这意味着Observable.subscribe() 不应使用Actor 的封闭范围来访问其内部状态。例如,您不应该在订阅中调用sender。相反,将其存储到一个 val 中,然后访问这个 val,如下例所示。
    • Akka 使用的调度器的分辨率与 Rx 不同。它的默认分辨率为 100 毫秒(Wyatt 2013)。如果有人遇到过这可能导致的问题,请在下面发表评论!

    最后,我为 Observables 实现了 ask 模式的等价物。它使用toObservable?? 异步返回一个Observable,由一个临时演员和一个PublishSubject 在幕后支持。请注意,源发送的消息是使用materialize()rx.lang.scala.Notification 类型,因此它们满足可观察合约中的completeerror 状态。否则我们无法将这些状态发送给接收器。但是,没有什么可以阻止您发送任意类型的消息;这些将简单地调用onNext()。如果在某个时间间隔内没有收到消息,则 observable 有一个超时并以超时异常停止。

    它是这样使用的:

    import akka.pattern.RX
    implicit val timeout = akka.util.Timeout(10 seconds)
    case object Req
    
    val system = ActorSystem("test")
    val source = system.actorOf(Props[Source],"thesource")
    
    class Source() extends Actor {
      def receive : Receive = {
         case Req =>
           val s = sender()
           Observable.interval(1 second).take(5).materialize.subscribe{s ! _}
      }
    }
    
    val obs = source ?? Req
    obs.observeOn(rx.lang.scala.schedulers.ExecutorScheduler(system.dispatcher)).subscribe((l : Any) => println ("onnext : " + l.toString),
                  (error : Throwable) => { error.printStackTrace ; system.shutdown() },
                  () => { println("completed, shutting system down"); system.shutdown() })
    

    并产生这个输出:

    onnext : 0
    onnext : 1
    onnext : 2
    onnext : 3
    onnext : 4
    completed, shutting system down
    

    来源如下。它是 AskSupport.scala 的修改版本。

    package akka.pattern
    
    /*
     * File : RxSupport.scala
     * This package is a modified version of 'AskSupport' to provide methods to 
     * support RX Observables.
     */
    
    import rx.lang.scala.{Observable,Subject,Notification}
    import java.util.concurrent.TimeoutException
    import akka.util.Timeout
    import akka.actor._
    import scala.concurrent.ExecutionContext
    import akka.util.Unsafe
    import scala.annotation.tailrec
    import akka.dispatch.sysmsg._
    
    class RxTimeoutException(message: String, cause: Throwable) extends TimeoutException(message) {
      def this(message: String) = this(message, null: Throwable)
      override def getCause(): Throwable = cause
    }
    
    trait RxSupport {
      implicit def toRx(actorRef : ActorRef) : RxActorRef = new RxActorRef(actorRef)
      def toObservable(actorRef : ActorRef, message : Any)(implicit timeout : Timeout) : Observable[Any] = actorRef ?? message
      implicit def toRx(actorSelection : ActorSelection) : RxActorSelection = new RxActorSelection(actorSelection)
      def toObservable(actorSelection : ActorSelection, message : Any)(implicit timeout : Timeout): Observable[Any] = actorSelection ?? message
    }
    
    final class RxActorRef(val actorRef : ActorRef) extends AnyVal {
      def toObservable(message : Any)(implicit timeout : Timeout) : Observable[Any] = actorRef match {
        case ref : InternalActorRef if ref.isTerminated =>
          actorRef ! message
          Observable.error(new RxTimeoutException(s"Recepient[$actorRef] has alrady been terminated."))
        case ref : InternalActorRef =>
          if (timeout.duration.length <= 0)
            Observable.error(new IllegalArgumentException(s"Timeout length must not be negative, message not sent to [$actorRef]"))
          else {
            val a = RxSubjectActorRef(ref.provider, timeout, targetName = actorRef.toString)
            actorRef.tell(message, a)
            a.result.doOnCompleted{a.stop}.timeout(timeout.duration)
          }
      }
      def ??(message :Any)(implicit timeout : Timeout) : Observable[Any] = toObservable(message)(timeout)
    }
    
    final class RxActorSelection(val actorSel : ActorSelection) extends AnyVal {
      def toObservable(message : Any)(implicit timeout : Timeout) : Observable[Any] = actorSel.anchor match {
        case ref : InternalActorRef =>
          if (timeout.duration.length <= 0)
            Observable.error(new IllegalArgumentException(s"Timeout length must not be negative, message not sent to [$actorSel]"))
          else {
            val a = RxSubjectActorRef(ref.provider, timeout, targetName = actorSel.toString)
            actorSel.tell(message, a)
             a.result.doOnCompleted{a.stop}.timeout(timeout.duration)
          }
        case _ => Observable.error(new IllegalArgumentException(s"Unsupported recipient ActorRef type, question not sent to [$actorSel]"))
      }
      def ??(message :Any)(implicit timeout : Timeout) : Observable[Any] = toObservable(message)(timeout)
    }
    
    
    private[akka] final class RxSubjectActorRef private (val provider : ActorRefProvider, val result: Subject[Any]) extends MinimalActorRef {
      import RxSubjectActorRef._
      import AbstractRxActorRef.stateOffset
      import AbstractRxActorRef.watchedByOffset
    
      /**
       * As an optimization for the common (local) case we only register this RxSubjectActorRef
       * with the provider when the `path` member is actually queried, which happens during
       * serialization (but also during a simple call to `toString`, `equals` or `hashCode`!).
       *
       * Defined states:
       * null                  => started, path not yet created
       * Registering           => currently creating temp path and registering it
       * path: ActorPath       => path is available and was registered
       * StoppedWithPath(path) => stopped, path available
       * Stopped               => stopped, path not yet created
       */
      @volatile
      private[this] var _stateDoNotCallMeDirectly: AnyRef = _
    
      @volatile
      private[this] var _watchedByDoNotCallMeDirectly: Set[ActorRef] = ActorCell.emptyActorRefSet
    
      @inline
      private[this] def watchedBy: Set[ActorRef] = Unsafe.instance.getObjectVolatile(this, watchedByOffset).asInstanceOf[Set[ActorRef]]
    
      @inline
      private[this] def updateWatchedBy(oldWatchedBy: Set[ActorRef], newWatchedBy: Set[ActorRef]): Boolean =
        Unsafe.instance.compareAndSwapObject(this, watchedByOffset, oldWatchedBy, newWatchedBy)
    
      @tailrec // Returns false if the subject is already completed
      private[this] final def addWatcher(watcher: ActorRef): Boolean = watchedBy match {
        case null => false
        case other => updateWatchedBy(other, other + watcher) || addWatcher(watcher)
      }
    
      @tailrec
      private[this] final def remWatcher(watcher: ActorRef): Unit = watchedBy match {
        case null => ()
        case other => if (!updateWatchedBy(other, other - watcher)) remWatcher(watcher)
      }
    
      @tailrec
      private[this] final def clearWatchers(): Set[ActorRef] = watchedBy match {
        case null => ActorCell.emptyActorRefSet
        case other => if (!updateWatchedBy(other, null)) clearWatchers() else other
      }
    
      @inline
      private[this] def state: AnyRef = Unsafe.instance.getObjectVolatile(this, stateOffset)
    
      @inline
      private[this] def updateState(oldState: AnyRef, newState: AnyRef): Boolean =
        Unsafe.instance.compareAndSwapObject(this, stateOffset, oldState, newState)
    
      @inline
      private[this] def setState(newState: AnyRef): Unit = Unsafe.instance.putObjectVolatile(this, stateOffset, newState)
    
      override def getParent: InternalActorRef = provider.tempContainer
    
      def internalCallingThreadExecutionContext: ExecutionContext =
        provider.guardian.underlying.systemImpl.internalCallingThreadExecutionContext
    
      /**
       * Contract of this method:
       * Must always return the same ActorPath, which must have
       * been registered if we haven't been stopped yet.
       */
      @tailrec
      def path: ActorPath = state match {
        case null =>
          if (updateState(null, Registering)) {
            var p: ActorPath = null
            try {
              p = provider.tempPath()
              provider.registerTempActor(this, p)
              p
            } finally { setState(p) }
          } else path
        case p: ActorPath       => p
        case StoppedWithPath(p) => p
        case Stopped =>
          // even if we are already stopped we still need to produce a proper path
          updateState(Stopped, StoppedWithPath(provider.tempPath()))
          path
        case Registering => path // spin until registration is completed
      }
    
      override def !(message: Any)(implicit sender: ActorRef = Actor.noSender): Unit = state match {
        case Stopped | _: StoppedWithPath => provider.deadLetters ! message
        case _ =>
          if (message == null) throw new InvalidMessageException("Message is null")
          else
            message match {
              case n : Notification[Any] => n.accept(result)
              case other                 => result.onNext(other)
            }
      }
    
      override def sendSystemMessage(message: SystemMessage): Unit = message match {
        case _: Terminate                      => stop()
        case DeathWatchNotification(a, ec, at) => this.!(Terminated(a)(existenceConfirmed = ec, addressTerminated = at))
        case Watch(watchee, watcher) =>
          if (watchee == this && watcher != this) {
            if (!addWatcher(watcher))
               // NEVER SEND THE SAME SYSTEM MESSAGE OBJECT TO TWO ACTORS
              watcher.sendSystemMessage(DeathWatchNotification(watchee, existenceConfirmed = true, addressTerminated = false))
          } else System.err.println("BUG: illegal Watch(%s,%s) for %s".format(watchee, watcher, this))
        case Unwatch(watchee, watcher) =>
          if (watchee == this && watcher != this) remWatcher(watcher)
          else System.err.println("BUG: illegal Unwatch(%s,%s) for %s".format(watchee, watcher, this))
        case _ =>
      }
    
      @deprecated("Use context.watch(actor) and receive Terminated(actor)", "2.2") override def isTerminated: Boolean = state match {
        case Stopped | _: StoppedWithPath => true
        case _                            => false
      }
    
      @tailrec
      override def stop(): Unit = {
        def ensureCompleted(): Unit = {
          result.onError(new ActorKilledException("Stopped"))
          val watchers = clearWatchers()
          if (!watchers.isEmpty) {
            watchers foreach { watcher =>
              // NEVER SEND THE SAME SYSTEM MESSAGE OBJECT TO TWO ACTORS
              watcher.asInstanceOf[InternalActorRef]
                .sendSystemMessage(DeathWatchNotification(watcher, existenceConfirmed = true, addressTerminated = false))
            }
          }
        }
        state match {
          case null => // if path was never queried nobody can possibly be watching us, so we don't have to publish termination either
            if (updateState(null, Stopped)) ensureCompleted() else stop()
          case p: ActorPath =>
            if (updateState(p, StoppedWithPath(p))) { try ensureCompleted() finally provider.unregisterTempActor(p) } else stop()
          case Stopped | _: StoppedWithPath => // already stopped
          case Registering                  => stop() // spin until registration is completed before stopping
        }
      }
    }
    
    private[akka] object RxSubjectActorRef {
      private case object Registering
      private case object Stopped
      private final case class StoppedWithPath(path : ActorPath)
    
      def apply(provider: ActorRefProvider, timeout: Timeout, targetName: String): RxSubjectActorRef = {
        val result = Subject[Any]()
        new RxSubjectActorRef(provider, result)
        /*timeout logic moved to RxActorRef/Sel*/
      }
    }
    /*
     * This doesn't work, need to create as a Java class for some reason ...
    final object AbstractRxActorRef {
        final val stateOffset = Unsafe.instance.objectFieldOffset(RxSubjectActorRef.getClass.getDeclaredField("_stateDoNotCallMeDirectly"))
        final val watchedByOffset = Unsafe.instance.objectFieldOffset(RxSubjectActorRef.getClass.getDeclaredField("_watchedByDoNotCallMeDirectly"))
    }*/
    
    package object RX extends RxSupport
    

    2015 年 9 月 10 日更新

    我想在这里添加一些更简单的代码来实现?? 运算符。这与上述略有不同,因为 a) 它不支持网络数据,b) 它返回 Observable[Observable[A]],这使得同步响应更容易。优点是它不会弄乱 Akka 内脏:

    object TypedAskSupport {
      import scala.concurrent.Future
      import akka.actor.{ActorRef,ActorSelection}
      import scala.reflect.ClassTag
    
      implicit class TypedAskableActorRef(actor : ActorRef) {
        val converted : akka.pattern.AskableActorRef = actor
        def ?[R](topic : Subscribe[R])(implicit timeout : akka.util.Timeout) : Future[Observable[R]] =
          converted.ask(topic).mapTo[Observable[R]]
        def ??[R](topic : Subscribe[R])(implicit timeout : akka.util.Timeout, execCtx : scala.concurrent.ExecutionContext) : Observable[Observable[R]] =
          Observable.from (this.?[R](topic)(timeout))
        def ?[R](topic : Request[R])(implicit timeout : akka.util.Timeout) : Future[R] =
          converted.ask(topic).asInstanceOf[Future[R]]
       def ??[R](topic : Request[R])(implicit timeout : akka.util.Timeout, execCtx : scala.concurrent.ExecutionContext) : Observable[R] =
          Observable.from { this.?[R](topic)(timeout) }
      }
    
      implicit class TypedAskableActorSelection(actor : ActorSelection) {
        val converted : akka.pattern.AskableActorSelection = actor
        def ?[R](topic : Subscribe[R])(implicit timeout : akka.util.Timeout) : Future[Observable[R]] =
          converted.ask(topic).mapTo[Observable[R]]
        def ??[R](topic : Subscribe[R])(implicit timeout : akka.util.Timeout, execCtx : scala.concurrent.ExecutionContext) : Observable[Observable[R]] =
          Observable.from (this.?[R](topic)(timeout))
        def ?[R](topic : Request[R])(implicit timeout : akka.util.Timeout) : Future[R] =
          converted.ask(topic).asInstanceOf[Future[R]]
      }
    }
    

    【讨论】:

      【解决方案3】:

      自从我发布原始问题以来,rx-java 和 akka 已经走了很长一段路。

      目前Akka Streams (middle of page) 有一个候选版本,我认为它在某种程度上试图提供与 rx-java 的Observable 类似的原语。

      还有一个针对Reactive Streams 的倡议,它看起来也通过toPublishertoSubscriber 方法提供了不同此类原语之间的互操作性; Akka 流实现了这个 API,java-rx 也有一个 extension 提供这个接口。在this blog post 上可以找到两者之间转换的示例,摘录如下:

      // create an observable from a simple list (this is in rxjava style)
      val first = Observable.from(text.split("\\s").toList.asJava);
      // convert the rxJava observable to a publisher
      val publisher = RxReactiveStreams.toPublisher(first);
      // based on the publisher create an akka source
      val source = PublisherSource(publisher);
      

      然后你大概可以在演员内部安全地传递这些。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2011-08-14
        • 2015-12-21
        • 1970-01-01
        • 1970-01-01
        • 2012-02-11
        • 1970-01-01
        • 1970-01-01
        • 2015-10-01
        相关资源
        最近更新 更多