【问题标题】:Is it possible to do these loops with parallel streams?是否可以使用并行流来执行这些循环?
【发布时间】:2020-05-12 01:37:07
【问题描述】:

如果我有在嵌套循环中执行的代码

int height = ...;
int width = ...;
int sourceIncrement = ...;
int destIncrement = ...;
int sourceOffset = ...;
int destOffset = ...;
int  sourceArray[] = ...;
byte destArray[] = ...;

for (int i = 0; i < height; i++) {                              
    for (int j = 0; j < width; j++) {               
        int pixel = sourceArray[sourceOffset ++];           
        destArray[destOffset ++] = (byte) (pixel      );
        destArray[destOffset ++] = (byte) (pixel >>  8);
        destArray[destOffset ++] = (byte) (pixel >> 16);
        destArray[destOffset ++] = (byte) (pixel >> 24);
    }                                           
    sourceOffset += sourceIncrement;                      
    destOffset += destIncrement;                     
}                                 

我想让它并行运行,我会尝试使用流来实现

IntStream.range(0, height).parallel().forEach(y -> {    
    IntStream.range(0, width).parallel().forEach(x -> {
        int pixel = sourceArray[sourceOffset ++];              
        destArray[destOffset ++] = (byte) (pixel      );   
        destArray[destOffset ++] = (byte) (pixel >>  8);   
        destArray[destOffset ++] = (byte) (pixel >> 16);   
        destArray[destOffset ++] = (byte) (pixel >> 24);   
    });                                            
    sourceOffset += sourceIncrement;                         
    destOffset += destIncrement;                        
}); 

但这是错误的,因为偏移量会在依赖于它们的代码运行之前增加。实际上是否有可能使其正确并行运行?

【问题讨论】:

  • 为内部循环的每次迭代计算 destOffset 和 sourceOffset。例如,它看起来只是sourceOffset = (i + sourceIncrement) * height + j。
  • @AndyTurner 那么我可以并行化外循环吗?内循环呢?
  • 如果您管理偏移量计算,则不需要两个循环。

标签: java loops optimization concurrency java-stream


【解决方案1】:

不要进行手动复制。

你可以使用,例如

IntBuffer src = IntBuffer.wrap(sourceArray, sourceOffset, sourceArray.length-sourceOffset);
IntBuffer dst = ByteBuffer.wrap(destArray, destOffset, destArray.length - destOffset)
        .order(ByteOrder.LITTLE_ENDIAN).asIntBuffer();

for(int i = 0; i < height; i++) {
    dst.put(src.limit(src.position()+width));
    src.limit(src.capacity()).position(src.position() + sourceIncrement);
    dst.position(dst.position() + (destIncrement >> 2));
}

这假设destIncrement 描述的是像素单位,即是四的倍数。此外,它假设数组足够长以至于有sourceIncrement resp。 destIncrement 最后一行的空间,没有写入,但是缓冲区不允许在最后一行之后设置位置。

如果数组没有那个空间,你必须在增加位置之前退出最后一行的循环:

if(height > 0) {
    IntBuffer src = IntBuffer.wrap(sourceArray, sourceOffset,
                                                sourceArray.length - sourceOffset);
    IntBuffer dst = ByteBuffer.wrap(destArray, destOffset, destArray.length - destOffset)
            .order(ByteOrder.LITTLE_ENDIAN).asIntBuffer();

    for(int i = 0; ;) {
        dst.put(src.limit(src.position()+width));
        if(++i == height) break;
        src.limit(src.capacity()).position(src.position() + sourceIncrement);
        dst.position(dst.position() + (destIncrement >> 2));
    }
}

目前尚不清楚这是否会从并行处理中受益,但为了完成,这里有一个并行变体:

IntBuffer src = IntBuffer.wrap(sourceArray, sourceOffset, (width+sourceIncrement)*height);
IntBuffer dst = ByteBuffer.wrap(destArray, destOffset, (width*4+destIncrement)*height)
        .order(ByteOrder.LITTLE_ENDIAN).asIntBuffer();

int srcRow = width + sourceIncrement, dstRow = width + (destIncrement>>2);
IntStream.range(0, height).parallel()
    .forEach(y -> dst.slice().position(y * dstRow)
        .put(src.slice().position(y * srcRow).limit(y * srcRow + width))
    );

由于缓冲区的位置和限制不是线程安全的,我们必须创建一个本地缓冲区。此解决方案使用slice(),它封装了sourceOffset 的初始值。 destOffset,简化转账操作的仓位和限位计算。在传输之前计算它们还可以确保该位置永远不会超过最后一行的width。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-03-29
    • 2015-11-10
    • 1970-01-01
    相关资源
    最近更新 更多