【问题标题】:Observable.Zip when number of sequences to zip is unknown until runtimeObservable.Zip 当要压缩的序列数在运行时未知时
【发布时间】:2012-09-22 01:01:48
【问题描述】:

我需要为审批流程建模。之前很简单。两个角色必须批准某件事,然后我们可以继续下一步:

public class Approved
{
    public string ApproverRole;
}

var approvals = Subscribe<Approved>();

var vpOfFinance = approvals.Where(e => e.ApproverRole == "Finance VP");
var vpOfSales = approvals.Where(e => e.ApproverRole == "Sales VP");

var approvedByAll = vpOfFinance.Zip(vpOfSales, Tuple.Create);

approvedByAll.Subscribe(_ => SomeInterestingBusinessProcess());

但现在有一个新要求:批准某事所需的角色数量可能会有所不同:

public class ApprovalRequested
{
    public string[] Roles;
}
var approvalRequest = Subscribe<ApprovalRequested>().Take(1);
var approvals = Subscribe<Approved>();

var approvedByAll = ???;

approvedByAll.Subscribe(_ => SomeInterestingBusinessProcess());

我觉得我在这里遗漏了一些非常明显的东西......谁能指出我正确的方向?

编辑

澄清一下:审批流程是基于每个项目的。批准到达的顺序是未定义的。我们不在乎一个角色是否多次批准一个项目。

【问题讨论】:

  • Zip 运算符希望成对的事情保持同步。在没有销售副总裁的情况下,您在这里所做的事情可能会得到财务副总裁的多次批准,并且事情可能会不同步。您需要在这里更好地定义您的要求。

标签: c# system.reactive observable


【解决方案1】:

问题基本上可以简化为从值流中创建Set,其中值可能是无序的或本质上很多。

如果 N 是集合的基数,我们可以简单地假设该过程将在至少 N 类型的值(在本例中为角色)被推送之前不会继续。

这是 Zip 运算符的示例解决方案;也许这可以让你开始:

    public static IObservable<IList<T>> Zip<T>(this IList<IObservable<T>> observables)
    {
        return Observable.Create<IList<T>>(observer =>
        {
            List<List<T>> store = new List<List<T>>(Enumerable.Range(1, observables.Count).Select(_ => new List<T>()));

            return new CompositeDisposable(observables.Select((o, i) => 
                o.Subscribe(value =>
                {
                    lock (store)
                    {
                        store[i].Add(value);

                        if (store.All(list => list.Count > 0))
                        {
                            observer.OnNext(store.Select(list => list[0]).ToList());
                            store.ForEach(list => list.RemoveAt(0));
                        }
                    }
                }))
            );
        });
    }

测试:

        Observable.Interval(TimeSpan.FromSeconds(0.5))
                  .GroupBy(i => i % 3)
                  .Select(gr => gr.AsObservable())
                  .Buffer(3)                      
                  .SelectMany(set => set.Zip())
                  .Subscribe(v => Console.WriteLine(String.Join(",", v)));

这里的一个问题是,在形成组时您可能会丢失初始值,因此您可能希望通过将方法重写为 IObservable&lt;IList&lt;T&gt;&gt; Zip&lt;TKey, T&gt;(this IGroupedObservable&lt;TKey, T&gt; observables) 来合并它。

【讨论】:

  • 我找不到这个 Zip 过载。我正在使用 nuget 运行 Rx 2.0.20823,包括 Rx-Experimental。
  • @JoãoBragança 我添加了一个 Zip 方法,它可以让您了解集合的基本概念。
  • Rx v2.0 中存在以下重载:静态 IObservable> Zip(此 IEnumerable> 来源)
【解决方案2】:

在当前版本的 Rx(我从 NuGet 获得)中,有一个版本的 Zip() 接受一个可观察对象的集合并返回一个可观察的集合。有了它,您可以执行以下操作:

string[] requiredApprovals = …;

var approvedByAll = requiredApprovals
    .Select(required => approvals.Where(a => a.ApproverRole == required))
    .Zip();

approvedByAll.Subscribe(_ => SomeInterestingBusinessProcess());

但正如@Enigmativity 所指出的,这只有在您可以确定每个人都以相同的顺序批准并且所有项目最终都会被所有必需角色批准的情况下才有效。如果没有,您将需要比Zip() 更复杂的东西。

【讨论】:

  • 请查看我对 Asti 回答的评论。
  • 我不知道为什么会这样。我也在使用 2.0.20823,它对我来说很好用。不过,我没有 Rx-Experimental。
  • 啊,我相信是因为你的目标是4.5,而我还在4.0。您对为什么不包含它有任何见解吗?
  • 我也想过这个问题,但是当面向 .Net 4.0 时它也适用于我(来自 VS 2012 和 2010)。
  • 我需要多加注意。 Zip 需要 IEnumerable>,我需要一些 IObservable>
猜你喜欢
  • 1970-01-01
  • 2012-12-07
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-11-23
  • 2021-11-27
相关资源
最近更新 更多