【问题标题】:ReplaySubject with distinct elements具有不同元素的 ReplaySubject
【发布时间】:2017-09-04 09:30:00
【问题描述】:

我想用 RxJava 实现一个 EventBus,我需要粘性事件。我知道我可以使用 BehaviorSubject,但它只缓存最后一个发出的项目,而我想缓存所有按其类型(类名)不同的事件。还有另一种选择 - ReplaySubject 但是它有一个开销 - 它包含所有发出的元素。 有没有办法创建某种类型的 ReplaySubject,它只包含唯一的类型元素?

【问题讨论】:

  • 您可以在将项目发送到Subject 之前对其进行过滤。因此,当它们在Subject 中收到时,它们在类型上已经是唯一的了。
  • @masp 这是一个事件总线。它就像一个无限流。我不知道我会发出多少物品,所以在将它们发送给主题之前我无法收集和过滤
  • 您想按类型缓存最新的元素吗?然后有多个 BehaviorSubjects 或 ReplaySubjects,每种类型一个,这也为您提供类型安全。
  • @akarnokd 不,我希望将所有元素发送到我的主题。但它们在类型上应该是唯一的。因此,如果我调用诸如 subject.onNext(a)、subject.onNext(b)、subject.onNext(a) 之类的东西,我希望缓存 [b, latest a]。我不需要类型安全,因为它是事件总线
  • 由于并发效应增加,这种情况需要具有非平凡内部结构的自定义主题。我宁愿完全避免使用事件总线,因为无论如何它们都比 ReactiveX 倒退了一步。

标签: rx-java2 behaviorsubject


【解决方案1】:

我不相信有任何干净的解决方案。不过,您或许可以完成这项工作。

  1. 为每个事件类型创建一个BehaviorSubject<>。使用ConcurrentMap 将每个传入事件分派到正确的Subject
  2. 按照事件到达的顺序维护这些主题的列表。
  3. 当接收到新订阅时,将创建一个可观察的对象,它是所有主题的合并,第一个主题列表已接收到事件,然后是未接收到事件的主题列表。

这是一些可以澄清上述内容的代码。未经测试。

// The subscription operation will perform a merge of the two lists
Map<EventType, BehaviorSubject<Event>> map = new ConcurrentHashMap<>();
List<BehaviorSubject<Event>> listOfUnseenEvents = new ArrayList<>();
List<BehaviorSubject<Event>> listOfSeenEvents = new ArrayList<>();
// ...
listOfUnseenEvents = map.values().asList();

public Observable<Event> busSub() {
  List<BehaviorSubject<Event>> allEvents = new ArrayList<>();
  synchronized ( map ) {
    allEvents.addAll( listOfSeenEvents );
    allEvents.addAll( listOfUnseenEvents );
  }
  return Observable.merge( allEvents );
}

// receive an event and dispatch it
eventSource
  .subscribe( event -> processEvent( event ) );

public void processEvent( Event event ) {
    BehaviorSubject<Event> eSubject = map.get( event.getEventType() );
    synchronized ( map ) {
      if ( containsEventType( listOfSeenEvents, event.getEventType() ) ) {
        removeEventType( listOfSeenEvents, event.getEventType() );
      } else {
        removeEventType( listOfUnseenEvents, event.getEventType() );
      }
      listOfSeenEvents.add( eSubject );
    }
    eSubject.onNext( event );
}

请注意,此代码利用了 merge() 将按顺序订阅每个给定的 observable 的事实。在 RxJava 文档中没有这样的保证。

如果事先不知道所有的事件类型,那么这是行不通的。

【讨论】:

  • 感谢您的想法,我不需要区分看到和未看到的事件,因此可以简化,如果所有事件类型都未确定,为什么这不起作用?我可以按事件类型对事件进行分组,似乎只发布特定类型的最新事件就足够了
  • 您需要“看不见的”事件主题,以便订阅者最终会看到这些类型的事件。否则,您需要为每个订阅者设置一个“未分类”主题,以及一个属于未分类的事件类型列表。
【解决方案2】:

我只用一个流无法成功解决这个问题。
所以我有一个按事件类型的流,并且只公开一个合并的流。

open class Event(val name: String)
class Event1(name: String) : Event(name)
class Event2(name: String) : Event(name)

val event1Subject = BehaviorSubject.create<Event1>()
val event2Subject = BehaviorSubject.create<Event2>()

event1Subject.onNext(Event1("event1 a"))
event1Subject.onNext(Event1("event1 b"))
event2Subject.onNext(Event2("event2 a"))
event2Subject.onNext(Event2("event2 b"))

Observable.merge(event1Subject, event2Subject)
        .subscribe { println(it.name) }

控制台输出:

event1 b
event2 b

【讨论】:

    猜你喜欢
    • 2019-03-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-09-28
    • 2020-03-04
    • 2013-09-14
    • 2018-12-11
    相关资源
    最近更新 更多