【问题标题】:Combining dependent observables结合依赖的可观察对象
【发布时间】:2016-02-17 23:42:34
【问题描述】:

我有一个容器类 Project,其中有很多 IItems。在项目启动之前,该项目可以添加项目。然后不能再添加任何项目,除非它被停止。

每个项目都可以通过其IsActive 属性依次激活和停用。

public interface IItem 
{
    bool IsActive { get; set;}
}

public interface IFoo : IItem
{
}

public class Foo : IFoo
{
    public bool IsActive { get; set;}
}

public interface IBar : IItem
{ }

public class Bar : IBar
{ 
    public bool IsActive { get; set;}
}

public class Project
{
    public Project(params IItem [] items)
    {
        Items = new List<IItem>(items);
    }

    public List<IItem> Items { get;}
}

我还有两个 observable,一个用于项目状态,一个用于对任何项目的更改。为了这个例子的目的,我已经用主题模拟了这些

var projectIsRunningObservable = new Subject<bool>();
var projectItemChangedObservable = new Subject<IItem>();

我正在尝试创建一个IObservable&lt;bool&gt;,它发送一个指示是否(至少一个项目处于活动状态并且项目已启动)的值。如果有活动项目并且项目已停止,它应该推送一个 false 值。

这是我目前所拥有的:

void Main()
{
    var bar1 = new Bar();   
    var bar2 = new Bar();   

    var foo1 = new Foo();
    var foo2 = new Foo();

    var projectIsRunningObservable = new Subject<bool>();
    var projectItemChangedObservable = new Subject<IItem>();

    var project = new Project(
        bar1, bar2, foo1, foo2);


    var observable = Observable.Create<bool>(obs =>
                {
                   IList<IItem> items = null;

                    var stateObservable = projectIsRunningObservable.StartWith(false).Subscribe(
                    (state) =>
                    {
                        if (!state)
                        {
                            items = null;
                            obs.OnNext(false);
                        }
                        else
                        {
                            items = project.Items.ToList();
                            obs.OnNext(items != null && items.Any(i => i.IsActive));
                        }
                    },
                    ex => obs.OnError(ex),
                    () => obs.OnCompleted());

                    var itemChangedObservable = projectItemChangedObservable.Subscribe(
                    x =>
                    {
                        obs.OnNext(items != null && items.Any(i => i.IsActive));
                    }
                    ,
                    ex => obs.OnError(ex),
                    () => obs.OnCompleted());

                    return new CompositeDisposable(stateObservable, itemChangedObservable);
                });


    var subscr = observable.Subscribe(Console.WriteLine);

    Console.WriteLine("Change bar1");
    bar1.IsActive = true;
    projectItemChangedObservable.OnNext(bar1);

    Console.WriteLine("Change bar2");
    bar2.IsActive = true;
    projectItemChangedObservable.OnNext(bar2);

    Console.WriteLine("Change foo1");
    foo1.IsActive = true;
    projectItemChangedObservable.OnNext(foo1);

    Console.WriteLine("Change foo2");
    foo2.IsActive = true;
    projectItemChangedObservable.OnNext(foo2);

    // Start project

    Console.WriteLine("Starting project");
    projectIsRunningObservable.OnNext(true);

    Console.WriteLine("Change bar1");
    bar1.IsActive = false;
    projectItemChangedObservable.OnNext(bar1);

    Console.WriteLine("Change bar2");
    bar2.IsActive = false;
    projectItemChangedObservable.OnNext(bar2);

    Console.WriteLine("Change foo1");
    foo1.IsActive = false;
    projectItemChangedObservable.OnNext(foo1);

    Console.WriteLine("Change foo2");
    foo2.IsActive = false;
    projectItemChangedObservable.OnNext(foo2);

    Console.WriteLine("Change foo2 back to true");
    foo2.IsActive = true;
    projectItemChangedObservable.OnNext(foo2);

    // Stop project
    Console.WriteLine("Stopping project");
    projectIsRunningObservable.OnNext(false);

    Console.WriteLine("Change bar1");
    bar1.IsActive = true;
    projectItemChangedObservable.OnNext(bar1);

    Console.WriteLine("Change bar2");
    bar2.IsActive = true;
    projectItemChangedObservable.OnNext(bar2);

    Console.WriteLine("Change foo1");
    foo1.IsActive = true;
    projectItemChangedObservable.OnNext(foo1);

    Console.WriteLine("Change foo2");
    foo2.IsActive = true;
    projectItemChangedObservable.OnNext(foo2);
}

这可行,但我不确定这是否是最好的方法,以及是否可以发送多个 OnErrorOnCompleted 通知。

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    这是Observable.Create 的一个很好的用例。您应该将Subject 重构为Observable,并使用TestScheduler 类进行测试。

    您的后一个问题的答案可以在Rx Design Guidelines (PDF) 中找到。第 6.2 章指出Observable.Create 提供了几种保护来使序列遵循 Rx 合同。

    当可观察序列完成时(通过触发 OnError 或 Oncompleted),任何订阅都将自动取消订阅。 任何订阅的观察者实例只会看到一个 OnError 或 OnCompleted 消息。不再发送消息。这样保证了OnNext*(OnError|OnCompleted)的Rx语法?

    注意:在指南示例中,他们使用Observable.CreateWithDisposable。在最新的 Rx 版本中,它已被重构为 Observable.Create 的重载,您可能知道 :)

    【讨论】:

    • 感谢您的回答,我不知道Observable.Create 会自动取消订阅OnErrorOnCompleted 上的任何订阅。小贴士!
    • @NedStoyanov - 所有可观察对象在完成(错误或其他情况)时应自动处理。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-01-25
    • 1970-01-01
    • 2013-08-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-05-19
    相关资源
    最近更新 更多