【发布时间】:2016-07-02 03:19:25
【问题描述】:
我有一个可观察的序列IObservable<int>,我想将其转换为IObservable<IList<int>>,同时保留以下要求:
- 最终的 observable 是一个批次序列,但是有两种——A和B,A批次包含1000每个项目,其中 B 个批次每个包含 400 个项目。
- 必须对原始序列中的每个数字进行两次批处理 - 一次在某个 A 批处理中,另一次在某个 B 批处理中
- 处理应在运行中进行,两种批次应并行生产。 IE。首先生产所有 A 批次,然后生产所有 B 批次的解决方案是不可接受的。
我可以使用Buffer 运算符轻松地生成一种批次,但我不知道如何在相同数据上生成两个批次。
编辑
这里有一个简单的代码来生成一个批次。
IObservable<int> source = GetSource(...);
await source
.Buffer(1000)
.Select(batch => Observable.FromAsync(() => ProcessBatchAsync(batch)))
.Merge(MaxConcurrentBatches)
.DefaultIfEmpty();
...
private async Task<Unit> ProcessBatchAsync(IList<int> batch)
{
...
return Unit.Default;
}
我想要的是:
- 要么是两个项目的 Observable,其中每个项目是仅一种批次的另一个 Observable。当批量可观察对象完成时,该主可观察对象应该是完整的。
- 一个可产生两种批次的 Observable。然后我需要根据 monad 的种类切换不同的运算符。
EDIT2
我需要详细说明约束。原始的 observable 位于 SqlReader 对象的顶部,订阅它两次意味着读取器被创建了两次,数据库访问量增加了一倍。我只需要一个订阅。
EDIT3
对于示例数据,我们可以使用Observable.Range(0,10000)。鉴于该顺序,我需要以任意顺序处理以下批次:
[0..1000), [0..400), [1000..2000),[400..800),[2000..3000),[800..1200),[3000..4000),[1200..1600) ... [9000..10000) ... [9600..10000)
或者您可以将范围 [0..100) 用于 10 和 4 个数字的批次。这并不重要,因为解决方案不应该取决于批次大小或批次类型的数量。
它应该适用于有 3 批 10、4 和 6 号,例如。或任何其他组合。
EDIT4
我认为我的限制让人们感到困惑。当我说“交错”时,我并不是说批次类型必须严格轮换。这正是我试图解释必须同时生产不同类型的批次的方式。给定 3 个批次类型 A、B 和 C,可能偶尔会一个接一个地生产两批 A 类型。但是,如果先是 A 型的所有批次,然后是 B 型的所有批次,然后是 C 型的所有批次,这是不可接受的。
【问题讨论】:
-
所以您希望列表中有 2 个订阅者具有不同的批量大小?
-
请展示您如何“使用
Buffer运算符轻松生产一种批次”,并为我们提供示例输入数据和预期输出。我可能很厚,但我真的明白你的解释。 -
@mark - 不要忘记
@通知 - 很幸运我刚好回来查看。 -
@mark - 您还没有得到问题中的“向我们提供示例输入数据和预期输出”。
-
@Enigmativity - 见 EDIT3。谢谢。
标签: c# system.reactive