标签: parallel-processing

"反序列化错误" - 带SOCK的foreach/doSNOW/snow(windows)

我正在使用SOCK集群和本地计算机上的工作程序运行并行操作.如果我限制我正在迭代的集合(在一次测试中使用70而不是完整的135个任务)那么一切正常.如果我去全套,我得到错误"反序列化错误(socklist [[n]]):从连接读取错误".

  • 我已取消阻止Windows防火墙中的端口(进/出)并允许Rscript/R的所有访问.

  • 它不能是超时问题,因为套接字超时设置为365天.

  • 它不是任何特定任务的问题,因为我可以顺序运行(如果我将数据集分成两半并进行两次单独的并行运行,也可以并行运行)

  • 我能想到的最好的是,通过套接字传输的数据太多了.似乎没有一个集群选项来限制数据限制.

我对如何进行感到茫然.有没有人见过这个问题或者可以建议修复?

这是我用来设置集群的代码:

cluster = makeCluster( degreeOfParallelism , type = "SOCK" , outfile = "" )
registerDoSNOW( cluster )
Run Code Online (Sandbox Code Playgroud)

编辑
虽然此问题与整个数据集有关,但它也会随着时间的推移而减少数据集.这可能表明这不仅仅是数据限制问题.

编辑2
我挖得更深一些,事实证明我的函数实际上有一个随机组件,这使得有时任务会引发错误.如果我按顺序运行任务,那么在操作结束时,我被告知哪个任务失败了.如果我并行运行,那么我会收到"unserialize"错误.我尝试在tryCatch调用中使用error = function(e){stop(e)}包装由每个任务执行的代码,但这也会生成"unserialize"错误.我很困惑因为我认为雪会把它们传回主人来处理错误?

parallel-processing foreach r

9
推荐指数
1
解决办法
3352
查看次数

PBS,刷新标准输出

我有一个长期运行的Torque/PBS工作,我想监控输出.但是只有在作业完成后才会复制日志文件.有没有办法说服PBS刷新它?

parallel-processing pbs batch-processing torque

9
推荐指数
2
解决办法
3431
查看次数

用于构造特里结构的并行算法?

因为trie数据结构具有如此巨大的分支因子,并且每个子树完全独立于其他子树,所以似乎应该有一种方法通过并行添加所有单词来极大地加速给定字典的构造.

我关于如何执行此操作的初步想法如下:将互斥锁与trie中的每个指针相关联(包括指向根的指针),然后让每个线程遵循用于将单词插入到trie中的常规算法.但是,在遵循任何指针之前,线程必须首先获取该指针的锁定,以便在需要向trie添加新的子节点时,它可以在不引入任何数据争用的情况下执行此操作.

这种方法的缺点是它使用了大量的锁 - 一个用于trie中的每个指针 - 并且执行大量的获取和释放 - 每个输入字符串中的每个字符一个.

有没有办法在没有使用几乎同样多的锁的情况下并行构建一个trie?

string algorithm parallel-processing trie data-structures

9
推荐指数
1
解决办法
1257
查看次数

C#parallel foreach同样完成任务

我正在使用C#Parallel.ForEach来处理超过数千个数据子集.一套需要5到30分钟来处理,具体取决于套装的大小.在我的电脑上有选项

ParallelOptions po = new ParallelOptions();
po.MaxDegreeOfParallelism = Environment.ProcessorCount
Run Code Online (Sandbox Code Playgroud)

我将获得8个并行进程.据我所知,流程在并行任务之间平均分配(例如,第一个任务获得工作号1,9,17等,第二个任务获得2,10,18等); 因此,一项任务可以比其他任务更快地完成自己的工作.因为这些数据集花费的时间比其他数据集少.

问题是四个并行任务在24小时内完成工作,但最后一个任务在48小时内完成.有没有机会组织并行性,以便所有并行任务完全平等?这意味着所有并行任务将继续有效,直到完成所有工作?

c# parallel-processing

9
推荐指数
1
解决办法
1122
查看次数

JUnit + Maven +并行测试执行错误

使用JUnit,Groovy,Spock和Maven时,我遇到了并行执行JUnit测试的问题.执行它们时,我在测试成功通过后得到以下内容:

[INFO] ------------------------------------------------------------------------
[INFO] BUILD FAILURE
[INFO] ------------------------------------------------------------------------
[INFO] Total time: 18.362s
[INFO] Finished at: Wed Mar 20 15:14:25 CET 2013
[INFO] Final Memory: 16M/221M
[INFO] ------------------------------------------------------------------------
[ERROR] Failed to execute goal org.apache.maven.plugins:maven-surefire-plugin:2.14:test (default-test) on project spock-webdriver: ExecutionException; nested exception is java.util.concurrent.ExecutionException: java.lang.RuntimeException: There was an error in the forked process
[ERROR] java.lang.NoSuchMethodError: org.apache.maven.surefire.common.junit4.JUnit4RunListener.rethrowAnyTestMechanismFailures(Lorg/junit/runner/Result;)V
[ERROR] at org.apache.maven.surefire.junit4.JUnit4Provider.invoke(JUnit4Provider.java:129)
[ERROR] at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
[ERROR] at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:57)
[ERROR] at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
[ERROR] at java.lang.reflect.Method.invoke(Method.java:601)
[ERROR] at org.apache.maven.surefire.util.ReflectionUtils.invokeMethodWithArray2(ReflectionUtils.java:208)
[ERROR] at org.apache.maven.surefire.booter.ProviderFactory$ProviderProxy.invoke(ProviderFactory.java:158)
[ERROR] at org.apache.maven.surefire.booter.ProviderFactory.invokeProvider(ProviderFactory.java:86)
[ERROR] …
Run Code Online (Sandbox Code Playgroud)

parallel-processing groovy junit maven spock

9
推荐指数
2
解决办法
6901
查看次数

R中的嵌套foreach循环更新公共数组

我试图在R中使用几个foreach循环来并行填充一个公共数组.我想要做的一个非常简化的版本是:

library(foreach)
set.seed(123)
x <- matrix(NA, nrow = 8, ncol = 2)

foreach(i=1:8) %dopar% {
    foreach(j=1:2) %do% {

      l <- runif(1, i, 100)
      x[i,j] <- i + j + l     #This is much more complicated in my real code.   

    }
}
Run Code Online (Sandbox Code Playgroud)

我想编码x并行更新矩阵,输出如下:

> x
       [,1]      [,2]
 [1,]  31.47017  82.04221
 [2,]  45.07974  92.53571
 [3,]  98.22533  12.41898
 [4,]  59.69813  95.67223
 [5,]  63.38633  55.37840
 [6,] 102.94233  56.61341
 [7,]  78.01407  69.25491
 [8,]  26.46907 100.78390 
Run Code Online (Sandbox Code Playgroud)

但是,我似乎无法弄清楚如何更新阵列.我试过把x <-它放到其他地方,但它似乎不喜欢它.我认为这将是一个非常容易解决的问题,但我所有的搜索还没有把我带到那里.谢谢.

parallel-processing foreach r

9
推荐指数
2
解决办法
8021
查看次数

等待网络会导致客户端超时吗?

我有一个服务器正在执行Azure队列指示的工作.它几乎总是在非常高的CPU上并行执行多个任务,并且一些任务使用Parallel.ForEach.在运行任务期间,我通过CloudQueue.AddMessageAsync使用await 调用将分析事件写入另一个Azure队列.

我注意到成千上万的这些分析文章因以下错误而失败:

WebException: The remote server returned an error: (500) Internal Server Error.

我检查了Azure的存储事件日志,我有一堆很好的PutMessage命令,端到端占用80.000ms,但它们只需要1ms用于Azure本身.我得到的HTTP状态代码是500,Azure描述了客户端超时的原因.

我认为正在发生的是我的代码调用AddMessageAsync并从那时起我的线程被释放,网络驱动程序正在发送请求并等待响应.获得响应时,网络驱动程序需要一个线程来获取响应,并且计划执行该任务并调用我的继续.由于我的服务器经常处于高负载状态,因此任务需要很长时间才能获得一个线程,然后Azure服务器会确定这是一个客户端超时.

调用azure的代码:

await cloudQueue.AddMessageAsync(new CloudQueueMessage(aMessageContent));
Run Code Online (Sandbox Code Playgroud)

例外:

StorageException: The remote server returned an error: (500) Internal Server Error.
Microsoft.WindowsAzure.Storage.Core.Executor.Executor.EndExecuteAsync[T](IAsyncResult result):11
Microsoft.WindowsAzure.Storage.Core.Util.AsyncExtensions+<>c__DisplayClass4.<CreateCallbackVoid>b__3(IAsyncResult ar):45
System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task):82
System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task):41
AzureCommon.Data.AsyncQueueDataContext+<AddMessage>d__d.MoveNext() in c:\BuildAgent\work\14078ab89161833\Azure\AzureCommon\Data\Async\AsyncQueueDataContext.cs:60
System.Runtime.CompilerServices.TaskAwaiter.ThrowForNonSuccess(Task task):82
System.Runtime.CompilerServices.TaskAwaiter.HandleNonSuccessAndDebuggerNotification(Task task):41
AzureCommon.Storage.AzureEvent+<DispatchAsync>d__1.MoveNext() in c:\BuildAgent\work\14078ab89161833\Azure\AzureCommon\Events\AzureEvent.cs:354

WebException: The remote server returned an error: (500) Internal Server Error.
System.Net.HttpWebRequest.EndGetResponse(IAsyncResult asyncResult):41
Microsoft.WindowsAzure.Storage.Core.Executor.Executor.EndGetResponse[T](IAsyncResult getResponseResult):44
Run Code Online (Sandbox Code Playgroud)

我说的为什么会这样吗?如果是这样,那么使用单线程同步上下文对我来说会更好吗?

Azure存储日志中的一行.您可以在此处找到有关每个属性含义的详细信息.

<request-start-time>            <operation-type>     <request-status> …
Run Code Online (Sandbox Code Playgroud)

c# parallel-processing async-await

9
推荐指数
1
解决办法
528
查看次数

gfortran是否利用DO CONCURRENT?

我目前正在使用gfortran 4.9.2,我想知道编译器是否真的知道如何利用DO CONCURRENT构造(Fortran 2008).我知道编译器"支持"它,但不清楚它是什么.例如,如果打开自动并行化(指定了一定数量的线程),编译器是否知道如何并行化并发循环?

编辑:正如评论中提到的,关于SO的前一个问题与我的非常相似,但它是从2012年开始的,只有最新版本的gfortran已经实现了现代Fortran的最新功能,所以我认为值得询问2015年编译器的当前状态.

parallel-processing fortran gfortran

9
推荐指数
1
解决办法
1006
查看次数

分配给"lib/ruby​​/2.1.0/timeout.rb"的1GB内存

我在循环中使用Twitter,Mongo和Parallel来检索和存储数据.

内存利用率达到1.5GB +

GC怎么不清洗这个?

更新: 这是有问题的脚本.

allocated memory by location
-----------------------------------
 973409328  /Users/jordan/.rvm/rubies/ruby-2.1.5/lib/ruby/2.1.0/timeout.rb:82
 359655091  /Users/jordan/.rvm/gems/ruby-2.1.5/gems/json-1.8.3/lib/json/common.rb:155
  34706221  /Users/jordan/.rvm/rubies/ruby-2.1.5/lib/ruby/2.1.0/openssl/buffering.rb:182
  31767589  /Users/jordan/.rvm/rubies/ruby-2.1.5/lib/ruby/2.1.0/net/http/response.rb:368
  22055648  /Users/jordan/.rvm/gems/ruby-2.1.5/gems/parallel-1.6.1/lib/parallel.rb:183
  12129637  /Users/jordan/.rvm/gems/ruby-2.1.5/gems/addressable-2.3.8/lib/addressable/uri.rb:525
  11115133  /Users/jordan/.rvm/rubies/ruby-2.1.5/lib/ruby/2.1.0/net/protocol.rb:172
  10609088  /Users/jordan/.rvm/gems/ruby-2.1.5/gems/addressable-2.3.8/lib/addressable/idna/pure.rb:177
   8333448  /Users/jordan/.rvm/gems/ruby-2.1.5/gems/twitter-5.15.0/lib/twitter/base.rb:152
   6041744  /Users/jordan/.rvm/gems/ruby-2.1.5/gems/thread_safe-0.3.5/lib/thread_safe/non_concurrent_cache_backend.rb:8
   4857232  /Users/jordan/.rvm/gems/ruby-2.1.5/gems/addressable-2.3.8/lib/addressable/uri.rb:1477
   4583920  /Users/jordan/.rvm/rubies/ruby-2.1.5/lib/ruby/2.1.0/monitor.rb:241
   4524872  /Users/jordan/.rvm/gems/ruby-2.1.5/gems/memoizable-0.4.2/lib/memoizable/method_builder.rb:117
   4282752  /Users/jordan/.rvm/gems/ruby-2.1.5/gems/twitter-5.15.0/lib/twitter/base.rb:151
   4200641  /Users/jordan/.rvm/gems/ruby-2.1.5/gems/mongo-2.1.1/lib/mongo/monitoring/command_log_subscriber.rb:104
   3283047  /Users/jordan/.rvm/rubies/ruby-2.1.5/lib/ruby/2.1.0/net/http/response.rb:61
   3150696  /Users/jordan/.rvm/gems/ruby-2.1.5/gems/mongo-2.1.1/lib/mongo/server/monitor.rb:125


allocated memory by gem
-----------------------------------
1084770550  ruby-2.1.5/lib
 359655091  json-1.8.3
  53016839  addressable-2.3.8
  22069048  parallel-1.6.1
  18422826  twitter-5.15.0
  10829988  mongo-2.1.1
   8908392  memoizable-0.4.2
   6041744  thread_safe-0.3.5
   4904294  faraday-0.9.2
   3839455  other
   3382080  naught-1.1.0
   2429320  bson-3.2.6
   1123917  rubygems
    320962  rollbar-2.4.0
    205097 …
Run Code Online (Sandbox Code Playgroud)

ruby parallel-processing memory-leaks mongodb

9
推荐指数
1
解决办法
578
查看次数

即时向Java 8并行Streams添加元素

目标是在Java 8流的帮助下处理连续的元素流.因此,在处理该流时,将元素添加到并行流的数据源中.

StreamsJavadoc在"无干扰"部分中描述了以下属性:

对于大多数数据源,防止干扰意味着确保在流管道的执行期间根本不修改数据源.值得注意的例外是其源是并发集合的流,这些集合专门用于处理并发修改.并发流源是Spliterator报告CONCURRENT特性的源.

这就是在我们的尝试中使用ConcurrentLinkedQueue的原因,它返回true

new ConcurrentLinkedQueue<Integer>().spliterator().hasCharacteristics(Spliterator.CONCURRENT)
Run Code Online (Sandbox Code Playgroud)

没有明确说明,在并行流中使用时不得修改数据源.

在我们的示例中,对于流中的每个元素,递增的计数器值被添加到队列中,该队列是流的数据源,直到计数器大于N.通过调用queue.stream(),一切正常,顺序执行:

import static org.junit.Assert.assertEquals;
import java.util.Queue;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.Stream;

public class StreamTest {
    public static void main(String[] args) {
        final int N = 10000;
        assertEquals(N, testSequential(N));
    }

    public static int testSequential(int N) {
        final AtomicInteger counter = new AtomicInteger(0);
        final AtomicInteger check = new AtomicInteger(0);
        final Queue<Integer> queue = new ConcurrentLinkedQueue<Integer>();

        for (int i = 0; i < N / 10; ++i) {
            queue.add(counter.incrementAndGet());
        }

        Stream<Integer> …
Run Code Online (Sandbox Code Playgroud)

java parallel-processing concurrency multithreading java-stream

9
推荐指数
1
解决办法
1053
查看次数