【问题标题】:Using Rx to Pair Up Work with Workers使用 Rx 与工作人员配对
【发布时间】:2015-12-19 23:33:49
【问题描述】:

作为 Rx 的新手,我正在尝试实现一些在概念上看起来很简单的东西,但我已经苦苦挣扎了一天,试图弄清楚如何在 Rx 中实现它。

我有 2 个 observable,一个是 IObservable<Worker>,另一个是 IObservable<Work>。这个想法是一种工作平衡器设置,其中工作通过一个 observable 提交,并与另一个 observable 上提交的 worker 配对。

使用Observable.Zip 很容易实现简单的情况。但是还有一个额外的要求在工作中引发了麻烦。

每个工作人员只能在有限的时间内接受工作(有一个窗口关闭,工作人员不再有资格接受工作,不应与工作项配对)。

这是我要实现的目标的大理石图:

------A-----B-----C------D-----> (WORKERS IObservable>) | | | | | | | | ------------]| ----] | | | |----------] | |-----] --------1---------2--3------> (WORK IObservable) --------A1---------C2---D3--> (WORKER/WORK PAIRS IObservable)

总结一下,要求:

  1. 当工作进入时,它应该分配给第一个工作进入时窗口仍处于活动状态的工作人员。(在图中我们有 A1 对但没有 B2,因为工作人员 B 已过期工作项 2 进入和 C2 成为正确配对的时间)

  2. 当一个worker与一个工作项配对时,它不应该有资格与后续的工作项配对(从某种意义上说,它的窗口应该在使用时强制关闭,但这是我的部分'我最挣扎如何实现)。

  3. 如果工作进入并且没有可用的工作人员打开窗口,它应该等到下一个工作人员发出并与之配对(参见图中的 D3)。

我很确定这可以通过Join 或GroupJoin 以某种方式实现,但我正在努力弄清楚如何获得我正在寻找的确切语义。任何帮助将不胜感激。

【问题讨论】:

  • 您已将其标记为 RxJava,但您的约定看起来让人想起 C#。这是用于 Java 还是 .NET?
  • 特定示例是用 C# 编写的,但我们同时使用这两种语言。我对如何使用通用 RX 结构在概念上解决这个问题更感兴趣,而不必担心特定于语言的语法。我认为 rx-java 和 .net 社区都可以做出贡献。
  • 对于工作人员的超时,您也许可以使用timeout 运算符,因为这符合语义。至于工作的配对,我可能建议不要试图将其硬塞到 Rx 运算符的组合中,因为这看起来逻辑很重。虽然如果练习是使用运算符,there are other combining operators 可能比Zip/GroupJoin 工作得更好(例如withLatestFrom)。
  • 为了从运行中移除一个忙碌的工作人员,它的 observable 可以在它被分配工作时完成(如果你想让它成为一个自定义的 observable)或者你可以有一个中间人(各种各样的)配对而跟踪谁在忙(而不是使用组合运算符)。

标签: system.reactive reactive-programming rx-java


【解决方案1】:

对于这项任务,我对join 的信任度不够,因为存在非确定性窗口和交叉评估。相反,我会为此编写一个自定义 zip 运算符,其中第一个序列的值带有时间限制值标记,并且 zip 算法只是丢弃第一个源的旧值。

Here is an example 类 ZipUntil 和 C# 中的可运行程序。

它是无锁的,并且使用了 Advanced RxJava 的概念,虽然我不明白同步取消在 Rx.NET 中是如何工作的(因为相同的模式在 RxJava 中不起作用)。

【讨论】:

  • 我真的很喜欢这个想法,尽管除非我要经常使用相同的模式,否则我只会实现一个不需要处理并发问题的更简单的版本(我可以运行一个后台线程并在同一线程上观察两个序列)
  • 这也给了我一个自定义 zip 运算符的想法,该运算符在 zip 时接受一个表达式来评估以丢弃与表达式不匹配的元素。类似Zip(IObs<T1>, IObs<T2>, Func<T1,bool>, Func<T2, bool>, Func<T1,T2,TResult> )
  • 我喜欢更好地控制我的操作符,是的,一旦理解了结构,就可以进行各种扩展。
【解决方案2】:

我怀疑这是最优的,但我认为它会起作用。

public static IObservable<Tuple<Worker, Work>> PairWorkers(IObservable<IObservable<Worker>> workerObservable, IObservable<Work> workObservable)
{
    var joined = workerObservable.Join(workObservable,
        worker => worker,
        work => Observable<Work>.Never(),
        (worker, work) => Tuple.Create(worker.FirstAsync(), work)
    );
    var distinct = joined.Distinct(t => t.Item2);
    return distinct;
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-07-01
    • 2021-08-15
    • 2020-06-14
    • 2019-04-22
    相关资源
    最近更新 更多