标签: parallel-processing

ListenableFuture回调执行顺序

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)

java parallel-processing guava java-7

5
推荐指数
1
解决办法
3658
查看次数

Python scikit学习n_jobs

这不是一个真正的问题,但我想了解:

  • 在Win7 4核8 GB系统上从Anaconda distrib运行sklearn
  • 在200.000个样本* 200个值表上拟合KMeans模型。
  • 使用n-jobs = -1运行(在将if __name__ == '__main__':行添加到脚本后),我看到脚本以4个进程开始,每个进程有10个线程。每个进程使用大约25%的CPU(总计:100%)。似乎按预期工作
  • 使用n-jobs = 1运行:停留在单个进程上(不足为奇),具有20个线程,并且还使用100%的CPU。

我的问题:如果库仍然使用所有内核,使用n-jobs(和joblib)有什么意义?我想念什么吗?它是Windows特定的行为吗?

python parallel-processing scikit-learn joblib

5
推荐指数
2
解决办法
1万
查看次数

如何在Python中正确使用多处理模块?

我有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)

python parallel-processing multithreading multiprocessing

5
推荐指数
1
解决办法
235
查看次数

Temporal多线程和超线程有什么区别?

有两个术语:

  • 时间多线程:在细粒度时间多线程中,主处理器流水线可以包含多个线程,其中上下文切换有效地发生在管道级之间(例如,在桶处理器中).桶处理器是在每个周期内在执行线程之间切换的CPU.

  • 超线程:是一种多线程,它允许单个处理器执行不同的线程,而无需同时真正执行它们.1这将其定性为时间切片或时间多线程而非同步多线程(SMT).观察到,由于长延迟事件,处理器的功能单元偶尔会在执行来自一个线程的指令时处于空闲状态.超线程试图通过执行来自另一个线程的指令来利用其他未使用的处理器周期,直到前一个线程准备好恢复执行.

TM和ST之间的主要区别在于,时间多线程(细粒度)使用C-slowing并在每个周期的执行线程之间切换,但是超线程在线程之间切换而不是每个周期并且仅在处理器的功能单元处于空闲状态时切换由于长延迟事件而从一个线程执行指令?

时间多线程(细粒度)和超线程之间有什么区别?

parallel-processing cpu concurrency multithreading cpu-architecture

5
推荐指数
1
解决办法
950
查看次数

在R中并行计算时更改核心数

我正在使用parallel包在R中执行并行化代码mclapply,并将预定义数量的内核作为参数.

如果我有一份将要运行几天的工作,我是否有办法编写(或换行)我的mclapply功能,以便在服务器高峰时段使用更少的内核,并在非高峰时段提高使用率?

parallel-processing r

5
推荐指数
1
解决办法
496
查看次数

Parallel Stream提供null项,如何在Java 8中完成

在尝试使用并行流时,我得到了一些奇怪的结果,我知道一种解决方法,但它似乎并不理想

// 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

selectedHashSet<Integer>在此操作期间未修改的.

somethingSubGroupDTO.addFilterDTO是一个辅助函数,它添加到一个不同步的List.这就是问题.作为一个非同步列表,我在列表中获得的项目少于预期,有些项目为空.如果我将其转换为同步列表,它可以工作.显然,将锁争用添加到并行流并不理想.

在高级别,我知道可以这样做,即每个流将自己进行处理,当它们加入时,它们将聚合.(至少我可以想象这样一个没有锁争用的进程)但是因为我是Java 8流处理的新手,所以我不知道怎么做.如何在没有单点争用的情况下执行相同的操作?

java parallel-processing lambda java-8 java-stream

5
推荐指数
1
解决办法
918
查看次数

即使我在父进程中使用sleep(),我的子进程也会在最后执行

有人可以解释为什么父进程总是在子进程中的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具有"更多"或独占权限的情况?

感谢帮助.:)

c parallel-processing sleep fork process

5
推荐指数
1
解决办法
1117
查看次数

是否存在允许并行加法的自然数的代数表示?

自然数可以使用二进制表示形式在对数空间中表示(此处为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

5
推荐指数
1
解决办法
170
查看次数

R:如何将包提供的方法导出到PSOCK群集?

R,我有一个矩阵列表,并希望将摘要函数应用于列表中的矩阵.矩阵代表社交网络,因此我需要应用ergm包提供的一些专门的汇总功能.这些摘要统计信息包含在摘要方法中.我可以编写一个函数作为这个汇总方法的包装器,并lapply用于将函数应用于矩阵列表.

但是,当我尝试通过使用parLapplyparSapplyparallel包中并行化时,结果看起来很奇怪.当我导出该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)

parallel-processing r

5
推荐指数
1
解决办法
2528
查看次数

我如何知道Java Stream收集(Collectors.toMap)是否已并行化?

我有以下代码尝试通过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)?

java parallel-processing java-stream

5
推荐指数
1
解决办法
1068
查看次数