【问题标题】:Is there an operator that works like Scan but let's me return an IObservable<TResult> instead of IObservable<TSource>?有没有像 Scan 一样工作的运算符,但让我返回一个 IObservable<TResult> 而不是 IObservable<TSource>?
【发布时间】:2016-07-21 08:28:58
【问题描述】:

在这个纯粹是为了练习而编造的例子中,这是我想要返回的内容:

如果两个学生在指定的时间段内加入学校,比如 2 秒,那么我想要一个数据结构返回这两个学生、他们加入的学校以及他们加入之间的时间间隔。

我一直在思考这些问题:

class Program
{
    static void Main(string[] args)
    {
        ObserveStudentsJoiningWithin(TimeSpan.FromSeconds(2));
    }

    static void ObserveStudentsJoiningWithin(TimeSpan timeSpan)
    {
        var school = new School("School 1");

        var admissionObservable =
            Observable.FromEventPattern<StudentAdmittedEventArgs>(school, "StudentAdmitted");

        var observable = admissionObservable.TimeInterval()
            .Scan((current, next) =>
            {
                if (next.Interval - current.Interval <= timeSpan)
                {
                    // But this won't work for me because
                    // this requires me to return a TSource
                    // and not a TResult
                }
            });

        var subscription = observable.Subscribe(TimeIntervalValueHandler);

        school.FillWithStudentsAsync(10, TimeSpan.FromSeconds(3));
        school.FillWithStudentsAsync(8, TimeSpan.FromSeconds(1));

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

这是域:

using System;
using System.Collections;
using System.Collections.Generic;
using System.Threading.Tasks;

namespace SchoolManagementSystem
{
    public class Student
    {
        private static int _studentNumber; 

        public Student(string name)
        {
            Name = name;
        }

        public string Name { get; set; }

        public static Student CreateRandom()
        {
            var name = string.Format($"Student {++_studentNumber}");

            return new Student(name);
        }

        public override string ToString()
        {
            return Name;
        }
    }

    public class School: IEnumerable<Student>
    {
        private List<Student> _students;

        public event StudentAdmitted StudentAdmitted;

        public string Name { get; set; }

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

        public void AdmitStudent(Student student)
        {
            if (!_students.Contains(student))
            {
                _students.Add(student);

                OnStudentAdmitted(this, student);
            }
        }

        protected virtual void OnStudentAdmitted(School school, Student student)
        {
            var args = new StudentAdmittedEventArgs(school, student);

            StudentAdmitted?.Invoke(this, args);
        }

        public IEnumerator<Student> GetEnumerator()
        {
            return _students.GetEnumerator();
        }

        IEnumerator IEnumerable.GetEnumerator()
        {
            return GetEnumerator();
        }
    }


    public delegate void StudentAdmitted(object sender, StudentAdmittedEventArgs args);

    public class StudentAdmittedEventArgs : EventArgs
    {
        public StudentAdmittedEventArgs(School school, Student student): base()
        {
            School = school;
            Student = student;
        }

        public School School { get; protected set; }
        public Student Student { get; protected set;  }
    }

    public static class Extensions
    {
        public async static void FillWithStudentsAsync(this School school, int howMany, TimeSpan gapBetweenEachAdmission)
        {
            if (school == null)
                throw new ArgumentNullException("school");

            if (howMany < 0)
                throw new ArgumentOutOfRangeException("howMany");

            if (howMany == 1)
            {
                school.AdmitStudent(Student.CreateRandom());
                return;
            }

            for (int i = 0; i < howMany; i++)
            {
                await Task.Delay(gapBetweenEachAdmission);

                school.AdmitStudent(Student.CreateRandom());
            }
        }
    }
}

但是,Scan 运算符允许我只返回相同 TSource 的 observable。 Select 也不会在这里工作,因为我不能向前看(我可以用Scan 做的事情)并将当前项目与下一个项目一起投影,即使Select 允许我将TSource 转换为TResult

我正在寻找介于两者之间的东西。

【问题讨论】:

  • Observable.Scan 确实允许您返回累加器类型。
  • 我的错。你说得对。如果你把它记下来作为答案?

标签: c# system.reactive


【解决方案1】:
  1. 对于成对比较(原始 - 将当前项目与下一个项目一起投影),您可以使用Buffer 方法来构建包含对的序列。
  2. 为了找出学生加入的间隔,使用Timestamp 而不是TimeInterval 方法可能更有用,因为下面的行next.Interval - current.Interval &lt;= timeSpan。你真正想要的是pair[1].Timestamp - pair[0].Timestamp &lt;= timeSpan

以下结果分为 4 对(学生 11,学生 12),(学生 13,学生 14),(学生 15,学生 16),(学生 17,学生 18):

var admissionObservable = Observable
        .FromEventPattern<StudentAdmittedEventArgs>(school, "StudentAdmitted")
        .Timestamp()
        .Buffer(2)
        .Where(pair => pair[1].Timestamp - pair[0].Timestamp <= timeSpan)
        .Select(pair => new JoiningData
        {
            Students = Tuple.Create(pair[0].Value.EventArgs.Student, pair[1].Value.EventArgs.Student),
            School = pair[0].Value.EventArgs.School,
            Interval = pair[1].Timestamp - pair[0].Timestamp
        });
  1. 正如@Enigmativity 所说,最好将每个元素与下一个元素进行比较。因此,为此我们可以使用Zip 方法:

以下结果为 8 对(学生 10,学生 11)(学生 11,学生 12),(学生 12,学生 13),(学生 13,学生 14),(学生 14,学生 15),(学生 15 , 学生 16), (学生 16, 学生 17), (学生 17, 学生 18):

var admissionObservable = Observable
     .FromEventPattern<StudentAdmittedEventArgs>(school, "StudentAdmitted")
     .Timestamp();        

admissionObservable
    .Zip(admissionObservable.Skip(1), (a, b) => Tuple.Create(a,b))        
    .Where(pair => pair.Item2.Timestamp - pair.Item1.Timestamp <= timeSpan)        
    .Select(pair => new JoiningData
    {
        Students = Tuple.Create(pair.Item1.Value.EventArgs.Student, pair.Item2.Value.EventArgs.Student),
        School = pair.Item1.Value.EventArgs.School,
        Interval = pair.Item2.Timestamp - pair.Item1.Timestamp
    });

【讨论】:

  • 确实很漂亮。非常感谢。我已经使用Scan 实现了它,并将很快发布。您对if 条件是正确的。我不得不改变它,因为它是错误的。但这是一个非常有趣的方法。因为在使用Scan 完成我的解决方案之前,我实际上想知道我是否应该记录DateTime 一个学生在Admitted 事件中被录取,并将该信息与EventArgs 一起发送。事实证明,我可以不使用 Scan,但这很有趣。
  • 此查询将学生分组为 (0, 1), (2, 3), (4, 5) 等成对,但不检查 1 &amp; 23 &amp; 4 等。
  • @Enigmativity - 是的,否则如果合适,一个元素将被存储两次,例如(学生 10,学生 11)和(学生 11,学生 12)中的“学生 11”。为这种方法添加了解决方案。
  • @AndriyTolstoy - 这是我认为需要的。如果学生 1 加入并在两秒钟内学生 2 加入,那么您会希望该对返回。还不清楚的是,如果 1、2 和 3 在三秒内全部加入,是否也需要返回 1 和 3 对?
【解决方案2】:

你能试试这个,看看它是否能满足你的需求?

IObservable<EventPattern<StudentAdmittedEventArgs>[]> observable =
    admissionObservable
        .Publish(pxs =>
            pxs
                .Window(pxs, x => Observable.Timer(timeSpan))
                .Select(ys => ys.Take(2)))
        .SelectMany(ys => ys.ToArray())
        .Where(ys => ys.Skip(1).Any());

【讨论】:

  • 因为窗口关闭边界是基于时间的,所以这不会返回所有在 2 秒内加入的学生,然后取第二个那个批次/窗口吗?
  • @WaterCoolrv2 - 是的,.Take(2) 就是这样做的 - 它只获取窗口中的第一对。它应该产生正确的结果。我在发布之前自己测试过。
【解决方案3】:

这是我所做的:

static void ObserveStudentsJoiningWithin(TimeSpan timeSpan)
{
    var school = new School("School 1");

    var admissionObservable =
        Observable.FromEventPattern<StudentAdmittedEventArgs>(school, "StudentAdmitted");

    var observable = admissionObservable.TimeInterval()
        .Scan<TimeInterval<EventPattern<StudentAdmittedEventArgs>>, StudentPair>(null, (previousPair, current) =>
        {
            Debug.Print(string.Format($"Student joined after {current.Interval.TotalSeconds} seconds, timeSpan = {timeSpan.TotalSeconds} seconds"));

            var pair = new StudentPair();

            if (previousPair == null)
            {
                pair.FirstStudent = null;
                pair.SecondStudent = current.Value.EventArgs.Student;
                pair.IntervalBetweenJoining = current.Interval;
                pair.School = current.Value.EventArgs.School;

                return pair;
            }

            if (current.Interval <= timeSpan)
            {
                pair.FirstStudent = previousPair.SecondStudent;
                pair.SecondStudent = current.Value.EventArgs.Student;
                pair.IntervalBetweenJoining = current.Interval;
                pair.School = current.Value.EventArgs.School;

                return pair;
            }
            else
            {
                return default(StudentPair);
            }
        })
        .Where(p => (p != default(StudentPair)) && (p.FirstStudent != null));

    var subscription = observable.Subscribe(StudentPairValueHandler);

    school.FillWithStudents(4, TimeSpan.FromSeconds(1));
    school.FillWithStudents(2, TimeSpan.FromSeconds(10));
    school.FillWithStudents(3, TimeSpan.FromSeconds(2));
    school.FillWithStudents(2, TimeSpan.FromSeconds(5));
    school.FillWithStudents(5, TimeSpan.FromSeconds(0.6));

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

static void StudentPairValueHandler(StudentPair pair)
{
    if (pair != null && pair.FirstStudent != null)
    {
        Console.WriteLine($"{pair.SecondStudent.Name} joined {pair.School.Name} {Math.Round(pair.IntervalBetweenJoining.TotalSeconds, 2)} seconds after {pair.FirstStudent.Name}.");
    }
}

...

public class StudentPair
{
    public Student FirstStudent;
    public Student SecondStudent;
    public School School;
    public TimeSpan IntervalBetweenJoining;
}


public static class Extensions
{
    public static void FillWithStudents(this School school, int howMany)
    {
        FillWithStudents(school, howMany, TimeSpan.Zero);
    }

    public static void FillWithStudents(this School school, int howMany, TimeSpan gapBetweenEachAdmission)
    {
        if (school == null)
            throw new ArgumentNullException("school");

        if (howMany < 0)
            throw new ArgumentOutOfRangeException("howMany");

        if (howMany == 1)
        {
            school.AdmitStudent(Student.CreateRandom());
            return;
        }

        for (int i = 0; i < howMany; i++)
        {
            Thread.Sleep((int)gapBetweenEachAdmission.TotalMilliseconds);

            school.AdmitStudent(Student.CreateRandom());
        }
    }

    public async static void FillWithStudentsAsync(this School school, int howMany, TimeSpan gapBetweenEachAdmission)
    {
        if (school == null)
            throw new ArgumentNullException("school");

        if (howMany < 0)
            throw new ArgumentOutOfRangeException("howMany");

        if (howMany == 1)
        {
            school.AdmitStudent(Student.CreateRandom());
            return;
        }

        for (int i = 0; i < howMany; i++)
        {
            await Task.Delay(gapBetweenEachAdmission);

            school.AdmitStudent(Student.CreateRandom());
        }
    }
}

【讨论】:

  • 我明白你在这里做什么,但它不是非常惯用的 Rx。查看“更好”设计的其他答案。
  • @LeeCampbell 谢谢。我意识到这一点,并且正在慢慢改进我的 Rx 思维方式。另一方面,我必须告诉你,我发现你的书是一个很好的学习资源。这是我很长时间以来读过的最好的书之一。它写得很好,总体上很漂亮。我从阅读中学到了很多东西。我还没有完整地阅读它,因为我一遍又一遍地重新阅读它的部分内容并慢慢练习很多次。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多