【问题标题】:How to request an item from IObservable?如何从 IObservable 请求项目?
【发布时间】:2016-04-27 14:51:13
【问题描述】:

原来的帖子包含一个问题,我设法解决了,引入了很多共享可变状态的问题。现在,我想知道是否可以以纯函数方式完成。


可以按特定顺序处理请求。

对于每个订单i有一个有效性E(i)

处理请求应遵循三个条件

  1. 在获取第一个请求和处理它之间应该没有延迟

  2. 处理某个请求和处理下一个请求之间应该没有延迟

  3. 当处理请求有多个顺序时,应选择效率最高的一个


具体例子:

对于整数的无限列表,将它们打印出来,这样,素数通常比非素数更早

排序的有效性与我们在队列中有素数但打印非素数的次数相反


我在C#(显然不是素数)中的第一个解决方案使用了一些具有由并发优先级队列表示的共享可变状态的类。这很丑陋,因为我必须手动为类订阅事件并取消订阅它们,检查队列在其他消费者处理之前是否被一个中间消费者用尽等等。

为了重构它,我选择了响应式扩展库,它似乎解决了状态问题。我知道在以下情况下我无法使用它:


source 函数不接受任何内容并返回 IObservable<Request>

process 函数接受 IObservable<Request> 并且不返回任何内容

我必须编写一个reorder 函数,它将请求从source 重新排序到process

reorder 内部有一个ConcurrentPriorityQueue 的订单。它应该处理两种情况:

  1. process 忙于处理reorder 找到更好的排序并更新队列

  2. process 请求新订单时,reorder 返回队列中的第一个元素


问题是如果reorder 返回IObservable<Request>,它不知道是否向它请求了项目,或者没有。

如果reorder 收到后立即致电OnNext,则它没有重新订购任何东西,违反了条件 3。

如果它确保找到了最佳排序,则它违反了条件 1 和 2,因为 process 可能会空闲。

如果reorder 返回ISubject<Request>,它会向消费者公开调用OnErrorOnCompleted 的选项。

如果reorder 已经返回队列,我会回到我开始的地方


问题是冷的 IObservable.Create 不够懒惰。当订阅它时,它开始耗尽所有请求的队列,但只使用了第一个请求的结果。

我想出的解决方案是返回可观察到的请求,即 IObservable<Func<Task<int>>> 而不是 IObservable<int>

它在只有一个订阅者时有效,但如果使用的请求多于源生成的数量,它们将永远等待。

这个问题可能可以通过引入缓存来解决,但是快速消费队列的消费者会对所有其他消费者产生副作用,因为他会以不太有效的顺序冻结队列,而不是等待一段时间后。

所以,我将发布原始问题的解决方案,但这并不是一个真正有价值的答案,因为它引入了很多问题。

这说明了为什么函数式反应式编程和副作用不能很好地融合在一起。另一方面,我现在似乎有一个无法以纯函数方式解决的实际问题的示例。还是我不?如果Order 函数接受optimizationLevel 作为参数,它将是纯粹的。我们能否以某种方式将时间隐式转换为 optimizationLevel 以使其也变得纯净?

我非常希望看到这样的解决方案。 C# 或任何其他语言。


有问题的解决方案。使用来自this repo 的 ConcurrentPriorityQueue。

using System;
using System.Collections.Generic;
using System.Threading.Tasks;
using System.Reactive.Linq;
using DataStructures;
using System.Threading;

namespace LazyObservable
{
    class Program
    {
        /// <summary>
        /// Compares tuple by second element, then by first in reverse
        /// </summary>
        class PriorityComparer<TElement, TPriority> : IComparer<Tuple<TElement, TPriority>>
            where TPriority : IComparable<TPriority>
        {
            Func<TElement, TElement, int> fallbackComparer;
            public PriorityComparer(IComparer<TElement> comparer=null)
            {
                if (comparer != null)
                {
                    fallbackComparer = comparer.Compare;
                }
                else if (typeof(IComparable<TElement>).IsAssignableFrom(typeof(TElement))
                    || typeof(IComparable).IsAssignableFrom(typeof(TElement)))
                {
                    fallbackComparer = (a,b)=>-Comparer<TElement>.Default.Compare(a,b);
                }
                else
                {
                    fallbackComparer = (_1,_2) => 0;
                }
            }
            public int Compare(Tuple<TElement, TPriority> x, Tuple<TElement, TPriority> y)
            {
                if (x == null && y == null)
                {
                    return 0;
                }
                if (x == null || y == null)
                {
                    return x == null ? -1 : 1;
                }
                int res=x.Item2.CompareTo(y.Item2);
                if (res == 0)
                {
                    res = fallbackComparer(x.Item1,y.Item1);
                }
                return res;
            }
        };
        const int N = 100;
        static IObservable<int> Source()
        {
            return Observable.Interval(TimeSpan.FromMilliseconds(1))
                .Select(x => (int)x)
                .Where(x => x <= 100);
        }
        static bool IsPrime(int x)
        {
            if (x <= 1)
            {
                return false;
            }
            if (x == 2)
            {
                return true;
            }
            int limit = ((int)Math.Sqrt(x)) + 1;
            for (int i = 2; i < limit; ++i)
            {
                if (x % i == 0)
                {
                    return false;
                }
            }
            return true;
        }
        static IObservable<Func<Task<int>>> Order(IObservable<int> numbers)
        {
            ConcurrentPriorityQueue<Tuple<int, int>> queue = new ConcurrentPriorityQueue<Tuple<int, int>>(new PriorityComparer<int, int>());
            numbers.Subscribe(x =>
            {
                queue.Add(new Tuple<int, int>(x, 0));
            });
            numbers
                .ForEachAsync(x=>
                {
                    Console.WriteLine("Testing {0}", x);
                    if (IsPrime(x))
                    {
                        if (queue.Remove(new Tuple<int, int>(x, 0)))
                        {
                            Console.WriteLine("Accelerated {0}", x);
                            queue.Add(new Tuple<int, int>(x, 1));
                        }
                    }
                });
            Func<Task<int>> requestElement = async () =>
              {
                  while (queue.Count == 0)
                  {
                      await Task.Delay(30);
                  }
                  return queue.Take().Item1;
              };
            return numbers.Select(_=>requestElement);
        }
        static void Process(IObservable<Func<Task<int>>> numbers)
        {
            numbers
                .Subscribe(async x=>
                {
                    await Task.Delay(1000);
                    Console.WriteLine(await x());
                });
        }

        static void Main(string[] args)
        {
            Console.WriteLine("init");
            Process(Order(Source()));
            //Process(Source());
            Console.WriteLine("called");
            Console.ReadLine();
        }
    }
}

【问题讨论】:

  • 我想,我几乎明白了。我现在正在研究冷的 IObservables。

标签: c# functional-programming system.reactive reactive-programming


【解决方案1】:

总结(概念上):

  1. 您有不规则传入的请求(来自source),并且有一个可以处理它们的处理器(函数process)。
  2. 处理器不应停机。
  3. 您隐含地需要某种队列式集合来管理请求进入速度超过处理器处理速度的情况。
  4. 如果有多个请求排队,理想情况下,您应该通过一些有效性函数对它们进行排序,但是重新排序不应成为停机的原因。 (函数reorder)。

这一切都正确吗?

假设是,source 可以是 IObservable&lt;Request&gt; 类型,听起来不错。 reorder 虽然听起来它真的应该返回一个IEnumerable&lt;Request&gt;process 想要在拉动的基础上工作:它想要在释放后拉取最高优先级的请求,如果队列中有则等待下一个请求为空,但立即开始。这听起来像是 IEnumerable 的任务,而不是 IObservable

public IObservable<Request> requester();
public IEnumerable<Request> reorder(IObservable<Request> requester);
public void process(IEnumerable<Request> requestEnumerable);

【讨论】:

  • 没错。但是当我创建一个IEnumerable&lt;Request&gt; 时,其中某处有一行yield return something();。万一,something()getValueWithHighestPriorityAfterCalculatingAllPriorities() 我会有停机时间。如果something()getFirstValueFromQueue() 又名getValueWithHighestPriorityBeforeCalculatingAllPriorities() 我没有重新排序。
  • 您可能必须将IEnumerable&lt;Request&gt; 实现为一个类。我在嘲笑一些东西。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-06-07
  • 2020-10-30
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-10-29
相关资源
最近更新 更多