【问题标题】:C# - Concurrent Foreach Like Thread SystemC# - 并发 Foreach 类线程系统
【发布时间】:2017-12-05 19:31:30
【问题描述】:
IObservable<Match> IObservableArray = new Regex("(.*):(.*)").Matches(file).OfType<Match>().ToList().ToObservable();
var query = IObservableArray.SelectMany(s => Observable.Start(() => {
    //do stuff
}));

上面的工作代码的解释:上面的代码使用 Observable 和 Reactive 来做一个并发多线程系统,同时保留 s 作为匹配。


我的问题是它似乎需要在开始执行 //do stuff 之前将所有内容加载到内存中,因为 IObservableArray 是一个很大的匹配数组 - 这会占用大量内存,导致它执行 OutOfMemory 异常。

我已经研究了一个多月,我能找到的只是 .Buffer() 如果我把它放在 .SelectMany() 之前,然后在 s 上进行 foreach 匹配,我能够将 1000 个匹配加载到内存中一段时间后,整体记忆力要好得多。

但是,由于我必须使用 foreach 一次遍历缓冲区中的所有 1000 个,所以它不是并发的——这意味着我基本上是一个接一个地检查 1。

有没有办法在下面做类似的代码,但它是并发/多线程的吗? (至少有 150 个并发运行,但不要全部加载到内存中,目前使用 1000 个。)

是的,我尝试使用 thread.start 等,使用它们可以更早地触发完成的代码,因为从技术上讲它确实完成了,因为它已经完成了被告知的事情,这使得它们全部进入一个新线程

IObservable<Match> IObservableArray = new Regex("(.*):(.*)").Matches(file).OfType<Match>().ToList().ToObservable();
var query = IObservableArray.Buffer(1000).SelectMany(s => Observable.Start(() => {
    //do stuff
}));
query.ObserveOn(ActiveForm).Subscribe(x =>
{
    //do finish stuff
});

【问题讨论】:

  • 无论如何使用 observables 背后的想法是什么?为什么不为此使用像 Parallel 这样的 TPL 类?
  • 这些实际上都不是多线程的,.ToList() 可能与 ToObersable() 重复,这可能会导致您消耗 2 倍内存
  • @PeterBons 老实说,我对此没有答案——我只是发现 Observables 有一种正确的方法来检测何时完成,所以就这样做了。
  • @user7842865 - 你能发一个minimal reproducible example吗?我们需要能够复制您的问题来解决它。
  • @user7842865 - 你能发一个minimal reproducible example吗?我们需要能够运行您的代码并复制您的问题。没有答案是因为你没有给我们minimal reproducible example

标签: c# multithreading foreach concurrency system.reactive


【解决方案1】:

您实际上并没有告诉Start() 要使用哪个调度程序,这可能是您没有获得所需并发性的原因。您可以将所需的调度程序指定为第二个参数:

var query = IObservableArray.Buffer(1000).SelectMany(s => Observable.Start(() => {
    //do stuff
}, TaskPoolScheduler.Default));

如果您知道 //do stuff 将花费超过 500 毫秒,我会考虑改用 ThreadPoolScheduler。在Task 阻塞线程至少 500 毫秒之前,任务池不会产生新线程,所以如果你知道你要做很多繁重的工作并且需要很多线程,你可以使用 @ 987654326@ 而不是 TaskPoolScheduler.Default

【讨论】:

    【解决方案2】:

    对于此类工作,IEnumerable&lt;T&gt;IObservable&lt;T&gt; 更适合。可枚举是您可以根据需要展开的东西,并在您准备好处理它们时获取它的值。相反,observable 是一种将其值强行推给您的东西,无论您是否能够处理负载。

    有多种方法可以并行处理IEnumerable&lt;T&gt;,并具有特定的并行度。在提出任何建议之前,要问的第一个问题是您必须与每个 Match 做的事情是同步的还是异步的。对于同步工作,最常用的工具是Parallel 类、PLINQTPL Dataflow 库。下面是一个 PLINQ 示例:

    IEnumerable<Match> matches = RegexFindAllMatches(file, "(.*):(.*)");
    Partitioner
        .Create(matches, EnumerablePartitionerOptions.NoBuffering)
        .AsParallel()
        .WithDegreeOfParallelism(Environment.ProcessorCount)
        .ForAll(match =>
        {
            // Do stuff
        });
    
    /// <summary>
    /// Provides an enumerable whose elements are the successful matches found by
    /// iteratively applying a regular expression pattern to the input string.
    /// </summary>
    public static IEnumerable<Match> RegexFindAllMatches(
        string input, string pattern, RegexOptions options = RegexOptions.None,
        TimeSpan matchTimeout = default)
    {
        if (matchTimeout == default) matchTimeout = Regex.InfiniteMatchTimeout;
        var match = Regex.Match(input, pattern, options, matchTimeout);
        while (match.Success)
        {
            yield return match;
            match = match.NextMatch();
        }
    }
    

    上述实现避免了使用Regex.Matches 方法,以及随后的MatchCollection 类,因为虽然这个类在枚举期间懒惰地评估下一个Match,然后它将每个找到的Match 存储在一个内部ArrayList (source code)。这可能会导致大量内存分配,与匹配的总数成正比。

    对于异步工作,Parallel 类和 PLINQ 不是很好的选择(除非您愿意等待 .NET 6 Parallel.ForEachAsync),但您仍然可以使用 TPL 数据流库。您还可以找到大量自定义选项herehere。搜索C# ForEachAsync 应该会显示更多选项。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-04-02
      • 1970-01-01
      • 2014-06-09
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多