我需要提升我的 python 应用程序。解决方案应该是微不足道的:
import time
from multiprocessing import Pool
class A:
def method1(self):
time.sleep(1)
print('method1')
return 'method1'
def method2(self):
time.sleep(1)
print('method2')
return 'method2'
def method3(self):
pool = Pool()
time1 = time.time()
res1 = pool.apply_async(self.method1, [])
res2 = pool.apply_async(self.method2, [])
res1 = res1.get()
res2 = res2.get()
time2 = time.time()
print('res1 = {0}'.format(res1))
print('res2 = {0}'.format(res2))
print('time = {0}'.format(time2 - time1))
a = A()
a.method3()
Run Code Online (Sandbox Code Playgroud)
但是每次我启动这个简单的程序时,我都会遇到一个异常:
Exception in thread Thread-2:
Traceback (most recent call last):
File "/usr/lib/python3.2/threading.py", line 740, in _bootstrap_inner
self.run() …Run Code Online (Sandbox Code Playgroud) 我有一个关于线程安全和互斥锁的问题。我有两个可能无法同时执行的函数,因为这可能会导致问题:
std::mutex mutex;
void A() {
std::lock_guard<std::mutex> lock(mutex);
//do something (should't be done while function B is executing)
}
T B() {
std::lock_guard<std::mutex> lock(mutex);
//do something (should't be done while function A is executing)
return something;
}
Run Code Online (Sandbox Code Playgroud)
现在的问题是,函数 A 和 B 不应该同时执行。这就是我使用互斥锁的原因。但是,如果从多个线程同时调用函数 B 是完全没问题的。但是,这也被互斥锁阻止了(我不想要这个)。现在,有没有办法确保 A 和 B 不会同时执行,同时仍然让函数 B 并行执行多次?
在Spring Batch的分区之间的关系gridSize的的PartitionHandler和数量的ExecutionContext通过传回的分区程序是有点混乱。例如,MultiResourcePartitioner声明它忽略 gridSize,但Partitioner文档没有解释何时/为什么可以接受。
例如,假设我有一个taskExecutor我想在不同的并行步骤中重复使用的对象,并且我将其大小设置为 20。如果我使用网格大小为 5的TaskExecutorPartitionerHandler,并且一个MultiResourcePartitioner返回任意数量的分区(每个文件一个),并行性实际上会如何表现?
假设MultiResourcePartitioner为特定运行返回 10 个分区。这是否意味着一次只执行其中的 5 个,直到所有 10 个都完成,并且这 20 个线程中不会有超过 5 个用于此步骤?
如果是这种情况,何时/为什么可以在Parititioner使用自定义实现覆盖时忽略 'gridSize' 参数?我认为如果在文档中对此进行了描述会有所帮助。
如果不是这种情况,我该如何实现?也就是说,我如何重新使用任务执行器并分别定义可以为该步骤并行运行的分区数量以及实际创建的分区数量?
好的,让我们开始吧,我脑子里有点混乱。
SEND:它正在阻塞。发送方将等待,直到接收方发布相应的 RECV。
SSEND:它是阻塞的,发送方不仅会等待接收方发布相应的 RECV,还会等待 RECV 的确认。这意味着 RECV 运行良好。
BSEND:它是非阻塞的。该进程可以继续执行其部分代码。数据存储在之前正确分配的缓冲区中。
ISEND:它是非阻塞的。该进程可以继续执行其部分代码。数据未存储在缓冲区中:在确定 ISEND 运行良好(WAIT/TEST)之前,您不得覆盖正在发送的数据。
那么.. ISEND 和 BSEND 仅在缓冲区上有所不同吗?
我想使用 Pandas 并行读取一个大的 .xls 文件。目前我正在使用这个:
LARGE_FILE = "LARGEFILE.xlsx"
CHUNKSIZE = 100000 # processing 100,000 rows at a time
def process_frame(df):
# process data frame
return len(df)
if __name__ == '__main__':
reader = pd.read_excel(LARGE_FILE, chunksize=CHUNKSIZE)
pool = mp.Pool(4) # use 4 processes
funclist = []
for df in reader:
# process each data frame
f = pool.apply_async(process_frame,[df])
funclist.append(f)
result = 0
for f in funclist:
result += f.get(timeout=10) # timeout in 10 seconds
Run Code Online (Sandbox Code Playgroud)
虽然这会运行,但我认为它实际上并没有加快读取文件的过程。有没有更有效的方法来实现这一目标?
我有 2 个测试套件。一个可以并行运行,另一个必须顺序运行。参见下面的定义。
我看到的是只有第二个运行。
我试图定义 2 个插件。没用。
我试图给他们不同的执行 ID。没用。
我试图将配置置于执行之下,但得到一个错误,该配置下的元素不被允许,例如failIfNoSpecifiedTests.
知道如何运行具有不同配置的套件 - 一个并行,另一个顺序?
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-surefire-plugin</artifactId>
<version>2.18.1</version>
<executions>
<execution>
<id>SequentialTests</id>
</execution>
</executions>
<configuration>
<includes>
<include>**/SequentialTests.java</include>
</includes>
</configuration>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-surefire-plugin</artifactId>
<version>2.18.1</version>
<executions>
<execution>
<id>ParallelTests</id>
</execution>
</executions>
<configuration>
<includes>
<include>**/ParallelTests.java</include>
</includes>
<threadCount>10</threadCount>
<parallel>classes</parallel>
</configuration>
</plugin>
Run Code Online (Sandbox Code Playgroud) 在哪些情况下哪个更有效?在某些情况下,哪一个根本无法工作?
我试图使一些通用代码更有效,并且很好奇哪个更好,因为据我所知,它们不能结合使用。
以供参考:
library(doParallel)
library(foreach)
foreach (i = list) %dopar% {
...
}
Run Code Online (Sandbox Code Playgroud)
对比
library(parallel)
parLapply(cl, X = list, fun = function)
Run Code Online (Sandbox Code Playgroud) 我正在寻找一种“整洁”且有效的方法来实现长步骤 1(可以并行化)和步骤 2 的组合,步骤 2 需要按原始顺序(如果可能的话,尽量减少来自第一步保存在 RAM 中)同时允许第二步在第一个对象的步骤 1 中的数据可用时立即开始,并与步骤 2 一起提供更多数据。
为了更详细地说明这一点,我需要压缩大量图像(慢速 - 第 1 步),然后通过网络连接按顺序发送每个图像(第 2 步)。在任何阶段限制 RAM 中准备好的压缩数据块的数量也很重要,例如,如果发送 1000 张图像,我想将“已完成”但未发送的图像数量限制为(例如)线程数/使用的处理器。
我已经完成了这个的“手写”版本,使用了一组 Task 对象,但它看起来很混乱,而且我相信其他人一定有类似的需求,所以有没有更“标准”的方法来做到这一点? 理想情况下,我希望有 2 个代表的 Parallel.ForEach 变体 - 一个用于第 1 步,一个用于第 2 步,我希望标准覆盖之一(例如包含“localFinal”参数的覆盖)可能有所帮助,但在原来这些最后阶段是“每个线程”,而不是“每个委托”。
任何人都可以指出我现有的巧妙方法来实现这一目标吗?
我必须运行很多随机森林模型,所以我想在我的 8 核服务器上使用 doParallel 来加速这个过程。
然而,某些模型需要比其他模型更长的时间,甚至可能会引发错误。我想并行运行 8 个模型,如果模型抛出错误和/或被跳过,那么工作人员应该继续。每个模型结果都保存在硬盘上,以便我以后可以访问和组合它们。
TryCatch
Run Code Online (Sandbox Code Playgroud)
或者
.errorhandling="remove"
Run Code Online (Sandbox Code Playgroud)
没有解决问题。我得到
Error in unserialize(socklist[[n]]) : error reading from connection
Run Code Online (Sandbox Code Playgroud)
代码示例:我用 %do% 试了一下,模型 2-7 运行成功。然而在 %dopar% 我得到了显示的错误
foreach(model=1:8, .errorhandling="remove") %dopar% {
tryCatch({
outl <- rf_perform(...)
saveRDS(outl,file=getwd() %+% "/temp/result_" %+% model %+% ".rds")
}, error = function(e) {print(e)}, finally = {})
}
Run Code Online (Sandbox Code Playgroud) 我参加了Parallel Programming 课程,它显示了并行接口:
def parallel[A, B](taskA: => A, taskB: => B): (A, B) = {
val ta = taskA
val tb = task {taskB}
(ta, tb.join())
}
Run Code Online (Sandbox Code Playgroud)
以下是错误的:
def parallel[A, B](taskA: => A, taskB: => B): (A, B) = {
val ta = taskB
val tb = task {taskB}.join()
(ta, tb)
}
Run Code Online (Sandbox Code Playgroud)
更多界面见https://gist.github.com/ChenZhongPu/fe389d30626626294306264a148bd2aa
它还向我们展示了执行四个任务的正确方法:
def parallel[A, B, C, D](taskA: => A, taskB: => B, taskC: => C, taskD: => D): (A, B, C, D) = {
val …Run Code Online (Sandbox Code Playgroud) python ×2
r ×2
c# ×1
c++ ×1
concurrency ×1
doparallel ×1
foreach ×1
java ×1
maven ×1
mpi ×1
mutex ×1
pandas ×1
performance ×1
scala ×1
spring-batch ×1