【发布时间】: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