【问题标题】:How to know which observable the next value comes from?如何知道下一个值来自哪个 observable?
【发布时间】:2016-07-20 23:34:11
【问题描述】:

假设我有这样的事情:

var observable = observable1
                 .Merge(observable2)
                 .Merge(observable3);

var subscription = observable.Subscribe(ValueHandler);

...

public void ValueHandler(string nextValue)
{
  Console.WriteLine($"Next value {nextValue} produced by /* sourceObservable */");
}

除了在每个 observable 实现中添加该引用以及值之外,有没有办法从产生下一个值的 observable1observable2observable3 中获取源 observable?

【问题讨论】:

  • 我也不明白你为什么这么快回答自己的问题。这是一个问答网站,而不是博客。
  • 不,不是很快。我真的在尝试,当我有问题时,我会发布它们,当我发布它们时,我同时仍在思考如何解决我的问题。碰巧我试图吸收我的想法,如果我认为我有答案,我会发布它。我不确定我的回答是否正确,我只想接受教育。就这样。除了学习之外,我真的没有发布任何议程。
  • 我的立场是正确的。
  • @Shlomo 没什么大不了的。 :-)

标签: c# system.reactive


【解决方案1】:

不,这与Merge 的设计目的完全相反。 Merge 旨在获取多个流并将它们视为一个。如果您想以某种方式将它们分开处理,请使用不同的运算符。


编辑

对于传递源的运算符,简短的回答是否定的。 Rx 是关于对消息做出反应,来源无关紧要。我什至不确定您是否可以从概念上定义 Rx 中的“源”是什么:

var observable1 = Observable.Interval(TimeSpan.FromSeconds(1)).Take(5);
var observable2 = observable1.Select(a => a.ToString());
var subscription = observable2.Subscribe(s => s.Dump());

订阅源是observable1observable2,还是某种指向系统时钟的指针?

如果您想将进入合并的消息分开,则可以使用Select,如下所示:

var observable = observable1.Select(o => Tuple.Create("observable1", o))
  .Merge(observable2.Select(o => Tuple.Create("observable2", o)))
  .Merge(observable3.Select(o => Tuple.Create("observable3", o)));

如果太乱,那么你可以很容易地制作一个扩展方法来清理它。


我还要补充一点,您在答案中发布的代码不是很像 Rx。一般准则是避免直接实施IObservableSchool 可以更简洁地改写如下:

    public class School
    {
        //private Subject<Student> _subject = null;
        private readonly ISubject<Student> _applicationStream = null;

        public static readonly int MaximumNumberOfSeats = 100;

        public string Name { get; set; }

        public School(string name)
            : this(name, new Subject<Student>())
        {
        }

        public School(string name, ISubject<Student> applicationStream )
        {
            Name = name;
            _applicationStream = applicationStream;
        }

        public void AdmitStudent(Student s)
        {
            _applicationStream.OnNext(s);
        }

        public IObservable<Student> ApplicationStream()
        {
            return _applicationStream;
        }

        public IObservable<Student> AcceptedStream()
        {
            return _applicationStream
                .SelectMany(s => s != null ? Observable.Return(s) : Observable.Throw<Student>(new ArgumentNullException("student")))
                .Distinct()
                .Take(MaximumNumberOfSeats);
        }
    }

通过这种方式,您可以订阅所有申请、接受以及拒绝等。您的状态也更少(没有List&lt;Student&gt;),理想情况下您甚至可以删除Subject&lt;Student&gt; applicationStream 和把它变成一个在某处传递的 Observable。

【讨论】:

  • 我认为问题不在于特定的运算符,而是没有一个运算符公开任何公共成员,这些成员提供有关生成值的上下文的任何附加信息。例如如果我使用Concat 而不是Merge,问题就会一直存在。
  • 不过,我希望得到纠正。是否有一个运算符可以结合两个或多个来源并告诉您哪个产生了值?
  • 没错。你是对的。 Rx 不关心其接口的实现。所以,这就是我所说的。你对Tuple sn-p 所做的就是改变TSourceIObservable&lt;TSource&gt;
  • 非常感谢您抽出宝贵时间。你的代码比我的更有意义,也更干净。实际上,我编写了我在答案中列出的代码的许多版本,其中一个与您所拥有的两个差异接近:我没有从外部收到 ISubject&lt;T&gt; 作为依赖项,而且我没有公开多个属性,例如 Accepted 等等,这些属性表示对 observable 的查询。我只是将学生暴露为可观察者。
  • 但是关于在学校班级内有状态(List&lt;T&gt;)是一个有争议的问题。我完全明白你的意思,有时我也倾向于使事情尽可能无国籍。纯可观察对象不需要有状态。但是在现实世界的场景中,如果没有状态,School 类会产生一些副作用(通过调用某个库将新录取的数据存储到数据库中)。我不会让它纯粹是无国籍的。但我不得不承认,我发现一个无状态的物体世界更加美丽和优雅。
【解决方案2】:

不,没有任何开箱即用的东西可以为我们提供有关哪个 observable 生成值的信息。

那是因为这是一个实现细节。

如果 Rx 的设计者必须提供这些附加信息,他们只能通过对 TSource 泛型类型参数施加某种限制来做到这一点。那可不好。这是值处理程序可以知道谁生成了值的唯一方法。

因此,获取此信息的责任在于开发人员使用 Rx 实现。

举个例子,假设您有一个School 课程,可以观察到这样的学生:

using System;
using System.Collections.Generic;
using System.Reactive.Subjects;


namespace SchoolManagementSystem
{
    public class Student
    {
        public Student(string name)
        {
            Name = name;
        }

        public string Name { get; set; }
    }

    public class School : IObservable<Student>
    {
        private List<Student> _students;
        private Subject<Student> _subject = null;

        public static readonly int MaximumNumberOfSeats = 100;

        public string Name { get; set; }

        public School(string name)
        {
            Name = name;
            _students = new List<Student>();
            _subject = new Subject<Student>();
        }

        public void AdmitStudent(Student student)
        {
            if (student == null)
            {
                var ex = new ArgumentNullException("student");
                _subject.OnError(ex);
                throw ex;
            }

            try
            {
                if (_students.Count == MaximumNumberOfSeats)
                {
                    _subject.OnCompleted();

                    return;
                }

                if (!_students.Contains(student))
                {    
                    _students.Add(student);

                    _subject.OnNext(student);
                }
            }
            catch(Exception ex)
            {
                _subject.OnError(ex);
            }
        }

        public IDisposable Subscribe(IObserver<Student> observer)
        {
            return _subject.Subscribe(observer);
        }
    }
}

客户端代码如下:

using SchoolManagementSystem;
using System;
using System.Reactive.Linq;

namespace Client
{
    class Program
    {
        static void Main(string[] args)
        {
            var school1 = new School("School 1");
            var school2 = new School("School 2");
            var school3 = new School("School 3");

            var observable = school1
                .Merge(school2)
                .Merge(school3);

            var subscription = observable
                .Subscribe(PrintStudentAdmittedMessage, PrintNoMoreStudentsCanBeAdmittedMessage);

            school1.FillWithStudents(100);
            school2.FillWithStudents(102);
            school3.FillWithStudents(101);

            Console.WriteLine("Press any key to stop observing and to exit the program.");
            Console.ReadKey();
            subscription.Dispose();
        }

        static void PrintStudentAdmittedMessage(Student student)
        {
            Console.WriteLine($"Student admitted: {student}");
        }

        static void PrintNoMoreStudentsCanBeAdmittedMessage()
        {
            Console.WriteLine("No more students can be admitted.");
        }
    }
}

那么,为了让客户知道一个学生被哪个学校录取,你必须改变IObservable&lt;TSource&gt;TSource的类型。在这种情况下,您必须更改 SchoolIObservable&lt;Student&gt; 的事实。

然而,改变它违背了领域模型的语义。因此,解决这个问题的方法是在学生内部包含对School 的引用,如下所示:

using System;
using System.Collections.Generic;

namespace SchoolManagementSystem
{
    public class Student
    {
        public Student(string name)
        {
            Name = name;
        }

        // Add this new property so you get this information
        // in the value handler
        public School School { get; set; }

        public string Name { get; set; }

        public override string ToString()
        {
            return string.Format($"({School.Name}: {Name})");
        }
    }
}

然后您将更改 School 类的 AdmitStudent 方法以指示学生被录取的学校,如下所示:

public void AdmitStudent(Student student)
{
    try
    {
        if (_students.Count == MaximumNumberOfSeats)
        {
            ...
        }

        if (!_students.Contains(student))
        {
            // Add this line to indicate which school
            // the student is being admitted to
            student.School = this;

            _students.Add(student);

            _subject.OnNext(student);
        }
    }
    catch(Exception ex)
    {
        _subject.OnError(ex);
    }
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-04-19
    • 2023-03-09
    • 2011-01-11
    • 1970-01-01
    • 2015-07-14
    • 2018-05-31
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多