【发布时间】:2016-04-27 14:51:13
【问题描述】:
原来的帖子包含一个问题,我设法解决了,引入了很多共享可变状态的问题。现在,我想知道是否可以以纯函数方式完成。
可以按特定顺序处理请求。
对于每个订单i有一个有效性E(i)
处理请求应遵循三个条件
在获取第一个请求和处理它之间应该没有延迟
处理某个请求和处理下一个请求之间应该没有延迟
当处理请求有多个顺序时,应选择效率最高的一个
具体例子:
对于整数的无限列表,将它们打印出来,这样,素数通常比非素数更早
排序的有效性与我们在队列中有素数但打印非素数的次数相反
我在C#(显然不是素数)中的第一个解决方案使用了一些具有由并发优先级队列表示的共享可变状态的类。这很丑陋,因为我必须手动为类订阅事件并取消订阅它们,检查队列在其他消费者处理之前是否被一个中间消费者用尽等等。
为了重构它,我选择了响应式扩展库,它似乎解决了状态问题。我知道在以下情况下我无法使用它:
source 函数不接受任何内容并返回 IObservable<Request>
process 函数接受 IObservable<Request> 并且不返回任何内容
我必须编写一个reorder 函数,它将请求从source 重新排序到process。
在reorder 内部有一个ConcurrentPriorityQueue 的订单。它应该处理两种情况:
当
process忙于处理reorder找到更好的排序并更新队列-
当
process请求新订单时,reorder返回队列中的第一个元素
问题是如果reorder 返回IObservable<Request>,它不知道是否向它请求了项目,或者没有。
如果reorder 收到后立即致电OnNext,则它没有重新订购任何东西,违反了条件 3。
如果它确保找到了最佳排序,则它违反了条件 1 和 2,因为 process 可能会空闲。
如果reorder 返回ISubject<Request>,它会向消费者公开调用OnError 和OnCompleted 的选项。
如果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