Guava的ListenableFuture库提供了一种将回调添加到未来任务的机制。这样做如下:
ListenableFuture<MyClass> future = myExecutor.submit(myCallable);
Futures.addCallback(future, new FutureCallback<MyClass>() {
@Override
public void onSuccess(@Nullable MyClass myClass) {
doSomething(myClass);
}
@Override
public void onFailure(Throwable t) {
printWarning(t);
}}, myCallbackExecutor);
}
Run Code Online (Sandbox Code Playgroud)
您可以ListenableFuture通过调用其get函数来等待a 完成。例如:
MyClass myClass = future.get();
Run Code Online (Sandbox Code Playgroud)
我的问题是,一定将来的所有回调都可以保证在get终止之前运行。即如果将来在许多回调执行程序上注册了许多回调,那么所有的回调会在get返回之前完成吗?
编辑
我的用例是,我将一个构建器传递给许多类。每个类填充构建器的一个字段。我希望所有字段都异步填充,因为每个字段都需要一个外部查询来生成该字段的数据。我希望给我打电话的用户asyncPopulateBuilder收到一个Future可以打电话的电话,get并确保所有字段都已填充。我认为的方法如下:
final Builder b;
ListenableFuture<MyClass> future = myExecutor.submit(myCallable);
Futures.addCallback(future, new FutureCallback<MyClass>() {
@Override
public void onSuccess(@Nullable MyClass myClass) {
b.setMyClass(myClass);
}
@Override
public void onFailure(Throwable t) {
printWarning(t); …Run Code Online (Sandbox Code Playgroud) 这不是一个真正的问题,但我想了解:
if __name__ == '__main__':行添加到脚本后),我看到脚本以4个进程开始,每个进程有10个线程。每个进程使用大约25%的CPU(总计:100%)。似乎按预期工作我的问题:如果库仍然使用所有内核,使用n-jobs(和joblib)有什么意义?我想念什么吗?它是Windows特定的行为吗?
我有110个PDF,我正在尝试从中提取图像.提取图像后,我想删除任何重复项并删除小于4KB的图像.我这样做的代码如下:
def extract_images_from_file(pdf_file):
file_name = os.path.splitext(os.path.basename(pdf_file))[0]
call(["pdfimages", "-png", pdf_file, file_name])
os.remove(pdf_file)
def dedup_images():
os.mkdir("unique_images")
md5_library = []
images = glob("*.png")
print "Deleting images smaller than 4KB and generating the MD5 hash values for all other images..."
for image in images:
if os.path.getsize(image) <= 4000:
os.remove(image)
else:
m = md5.new()
image_data = list(Image.open(image).getdata())
image_string = "".join(["".join([str(tpl[0]), str(tpl[1]), str(tpl[2])]) for tpl in image_data])
m.update(image_string)
md5_library.append([image, m.digest()])
headers = ['image_file', 'md5']
dat = pd.DataFrame(md5_library, columns=headers).sort(['md5'])
dat.drop_duplicates(subset="md5", inplace=True)
print "Extracting the unique images."
unique_images …Run Code Online (Sandbox Code Playgroud) 有两个术语:
时间多线程:在细粒度时间多线程中,主处理器流水线可以包含多个线程,其中上下文切换有效地发生在管道级之间(例如,在桶处理器中).桶处理器是在每个周期内在执行线程之间切换的CPU.
超线程:是一种多线程,它允许单个处理器执行不同的线程,而无需同时真正执行它们.1这将其定性为时间切片或时间多线程而非同步多线程(SMT).观察到,由于长延迟事件,处理器的功能单元偶尔会在执行来自一个线程的指令时处于空闲状态.超线程试图通过执行来自另一个线程的指令来利用其他未使用的处理器周期,直到前一个线程准备好恢复执行.
TM和ST之间的主要区别在于,时间多线程(细粒度)使用C-slowing并在每个周期的执行线程之间切换,但是超线程在线程之间切换而不是每个周期并且仅在处理器的功能单元处于空闲状态时切换由于长延迟事件而从一个线程执行指令?
时间多线程(细粒度)和超线程之间有什么区别?
parallel-processing cpu concurrency multithreading cpu-architecture
我正在使用parallel包在R中执行并行化代码mclapply,并将预定义数量的内核作为参数.
如果我有一份将要运行几天的工作,我是否有办法编写(或换行)我的mclapply功能,以便在服务器高峰时段使用更少的内核,并在非高峰时段提高使用率?
在尝试使用并行流时,我得到了一些奇怪的结果,我知道一种解决方法,但它似乎并不理想
// Create the set "selected"
somethingDao.getSomethingList().parallelStream()
.filter(something -> !selected.contains(something.getSomethingId()))
.forEach(something ->
somethingSubGroupDTO.addFilterDTO(
new FilterDTO(something.getSomethingName(), something.getSomethingDescription(), false))
);
selected.clear();
Run Code Online (Sandbox Code Playgroud)
somethingDao.getSomethingList 返回一个 List
selected是HashSet<Integer>在此操作期间未修改的.
somethingSubGroupDTO.addFilterDTO是一个辅助函数,它添加到一个不同步的List.这就是问题.作为一个非同步列表,我在列表中获得的项目少于预期,有些项目为空.如果我将其转换为同步列表,它可以工作.显然,将锁争用添加到并行流并不理想.
在高级别,我知道可以这样做,即每个流将自己进行处理,当它们加入时,它们将聚合.(至少我可以想象这样一个没有锁争用的进程)但是因为我是Java 8流处理的新手,所以我不知道怎么做.如何在没有单点争用的情况下执行相同的操作?
有人可以解释为什么父进程总是在子进程中的while循环开始之前完全完成,即使我让父进程休眠.
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/types.h>
#include <unistd.h>
int main(int argc, char **argv) {
int i = 10000, pid = fork();
if (pid == 0) {
while(i > 50) {
if(i%100==0) {
sleep(20);
}
printf("Child: %d\n", i);
i--;
}
} else {
while(i < 15000) {
if(i%50==0) {
sleep(50);
}
printf("Parent: %d\n", i);
i++;
}
}
exit(EXIT_SUCCESS);
}
Run Code Online (Sandbox Code Playgroud)
输出如下所示:
Parent: ..
Parent: ..
Parent: ..
Run Code Online (Sandbox Code Playgroud)
直到Parent完成,然后对于Child.
原因可能是我在单核CPU上进行测试?如果我有多核设置,结果会改变吗?sleep(50)肯定有效 - 因为脚本完成需要很长时间 - 为什么CPU不会切换进程?是否存在例如在while循环期间进程对CPU具有"更多"或独占权限的情况?
感谢帮助.:)
自然数可以使用二进制表示形式在对数空间中表示(此处为little-endian):
-- The type of binary numbers; little-endian; O = Zero, I = One
data Bin = O Bin | I Bin | End
Run Code Online (Sandbox Code Playgroud)
然后可以通过调用函数()次数来实现a和添加.这种实现的问题在于它本质上是顺序的.为了添加2个号码,呼叫按顺序链接.(例如使用进位)的其他实现遭受相同的问题.很容易看出,添加不能与该表示并行实现.是否使用代数数据类型表示自然数,它采用对数空间,并且可以并行添加?bsuccessorO(log(N))absucadd
插图代码:
-- The usual fold
fold :: Bin -> (t -> t) -> (t -> t) -> t -> t
fold (O bin) zero one end = zero (fold bin zero one end)
fold (I bin) zero one end = one (fold bin …Run Code Online (Sandbox Code Playgroud) algorithm parallel-processing haskell addition data-structures
在R,我有一个矩阵列表,并希望将摘要函数应用于列表中的矩阵.矩阵代表社交网络,因此我需要应用ergm包提供的一些专门的汇总功能.这些摘要统计信息包含在摘要方法中.我可以编写一个函数作为这个汇总方法的包装器,并lapply用于将函数应用于矩阵列表.
但是,当我尝试通过使用parLapply或parSapply从parallel包中并行化时,结果看起来很奇怪.当我导出该summary.statistics函数时,我甚至会收到一条错误消息.
我是否必须将ergm包提供的摘要方法导出到群集对象?如果是这样,怎么样?以下代码是一个独立的示例.
library("ergm")
library("parallel")
# create list of matrices
m <- matrix(rbinom(900, 1, 0.1), nrow = 30)
l <- list(m, m, m, m, m)
# write wrapper function that computes results
fun <- function(mat) {
s <- summary(mat ~ edges + dsp(1))
return(s)
}
cl <- makePSOCKcluster(2) # create cluster object
test1 <- sapply(l, fun) # works!
test2 <- parSapply(cl, l, fun) # problem: results …Run Code Online (Sandbox Code Playgroud) 我有以下代码尝试通过Java Stream API以并行方式从List填充Map:
class NameId {...}
public class TestStream
{
static public void main(String[] args)
{
List<NameId > niList = new ArrayList<>();
niList.add(new NameId ("Alice", "123456"));
niList.add(new NameId ("Bob", "223456"));
niList.add(new NameId ("Carl", "323456"));
Stream<NameId> niStream = niList.parallelStream();
Map<String, String> niMap = niStream.collect(Collectors.toMap(NameId::getName, NameId::getId));
}
}
Run Code Online (Sandbox Code Playgroud)
我如何知道是否使用多个线程填充地图,即并行?我是否需要调用Collectors.toConcurrentMap而不是Collectors.toMap?这是一种合理的方式来并行化地图的人口吗?我怎么知道具体的地图支持新的niMap(例如它是HashMap)?