【问题标题】:How to produce several batches of different size for the same data using Rx.NET?如何使用 Rx.NET 为相同的数据生成多个不同大小的批次?
【发布时间】:2016-07-02 03:19:25
【问题描述】:

我有一个可观察的序列IObservable<int>,我想将其转换为IObservable<IList<int>>,同时保留以下要求:

  • 最终的 observable 是一个批次序列,但是有两种——ABA批次包含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


【解决方案1】:

我认为这会产生你需要的东西:

var query =
    Observable
        .Range(0, 10000)
        .Publish(ns =>
            ns
                .Buffer(1000)
                .Concat(Observable.Repeat(new List<int>() as IList<int>))
                .Zip(ns.Buffer(400), (n1s, n2s) => new [] { n1s, n2s })
                .SelectMany(nns => nns)
                .Where(xs => xs.Any()));

根据您的示例输出正确交错。

如果我将数字减少 100 倍,那么我会得到以下输出:

0、1、2、3、4、5、6、7、8、9 0, 1, 2, 3 10、11、12、13、14、15、16、17、18、19 4、5、6、7 20、21、22、23、24、25、26、27、28、29 8、9、10、11 30, 31, 32, 33, 34, 35, 36, 37, 38, 39 12、13、14、15 40, 41, 42, 43, 44, 45, 46, 47, 48, 49 16、17、18、19 50, 51, 52, 53, 54, 55, 56, 57, 58, 59 20、21、22、23 60, 61, 62, 63, 64, 65, 66, 67, 68, 69 24、25、26、27 70, 71, 72, 73, 74, 75, 76, 77, 78, 79 28、29、30、31 80, 81, 82, 83, 84, 85, 86, 87, 88, 89 32、33、34、35 90, 91, 92, 93, 94, 95, 96, 97, 98, 99 36、37、38、39 40, 41, 42, 43 44、45、46、47 48、49、50、51 52、53、54、55 56、57、58、59 60, 61, 62, 63 64、65、66、67 68、69、70、71 72、73、74、75 76, 77, 78, 79 80、81、82、83 84、85、86、87 88、89、90、91 92, 93, 94, 95 96, 97, 98, 99

如果您不需要它们严格交错,那么这是一种概括 n 缓冲区的方法:

var buffers = new [] { 1000, 400, 500, 300 };
var source = Observable.Range(0, 10000);
var result = source.Publish(ss => buffers.Select(b => ss.Buffer(b)).Merge());

正如 Theo 指出的那样,第一个查询中使用的Observable.Repeat 有可能产生大量对象。如果使用Observable.Range,则不会发生这种情况,但如果使用Observable.Interval,则这是一个大问题。

我很懒惰,根本没有尝试通过计算来限制所需填充物的数量。

很容易修复。

var total_item_count = 100;
var batch_a_size = 10;
var batch_b_size = 4;

var filler =
    Observable
        .Repeat(
            new List<long>() as IList<long>,
            Math.Abs(total_item_count / batch_a_size - total_item_count / batch_b_size) + 1);

var query =
    Observable
        .Interval(TimeSpan.FromSeconds(0.01))
        .Take(total_item_count)
        .Publish(ns =>
            ns
                .Buffer(batch_a_size)
                .Concat(filler)
                .Zip(
                    ns
                        .Buffer(batch_b_size)
                        .Concat(filler),
                    (n1s, n2s) => new [] { n1s, n2s })
                .SelectMany(nns => nns)
                .Where(xs => xs.Any()));

我得到了和以前一样的输出,现在它可以针对每个批量大小进行配置,并且没有逃跑Repeat

【讨论】:

  • +1 获取创意解决方案。但是你能把它扩展到 N 批类型吗? ionoy 提供的解决方案更直接,可以轻松推广到 N 个批次类型。
  • @mark - Ionoy 的解决方案不会交错。输出不需要交错吗?
  • 我需要同时生成它们,换句话说,我确实需要将它们交错。原因 - 不同的消费者同时消费不同的批次。但是你确定 Ionoy 的解决方案不会交错吗? Subscribe 方法不会立即返回吗?我需要检查一下。
  • @TheodorZoulias - 确实如此。我已经解决了这个问题。
  • 现在它似乎在控制之中。令人恐惧的是,Rx 允许如此轻松地创建难以察觉的内存泄漏。
【解决方案2】:

下面的方法怎么样

var observable = Observable.Range(0, 10000);

Task.Run(() => observable.Buffer(400).Subscribe(buffer => /* process buffer */ ));
Task.Run(() => observable.Buffer(1000).Subscribe(buffer => /* process buffer */ ));

两种批次将并行生产。

【讨论】:

  • 请看我的 EDIT2。你的提议对我来说不是很好,因为这意味着我将对数据库的访问增加一倍。
  • 使用Publish(...) 创建一个包含多个订阅者的数据库订阅。
  • @ionoy - 请安排您的评论作为代码示例的答案。
【解决方案3】:

此代码创建两个IObservable&lt;IList&lt;int&gt;&gt; 序列。

var source = Observable.Range(0,10000)
                       .Publish();

var batchesOf1000 = source.Buffer(1000);
var batchesOf400 = source.Buffer(400);

batchesOf1000.Subscribe(batch => batch.Dump());
batchesOf400.Subscribe(batch => batch.Dump());

var disposable = source.Connect();

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-10-09
    • 2019-01-30
    • 2018-03-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-03-24
    • 2014-01-01
    相关资源
    最近更新 更多