You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

带依赖的连续cudaMemcpyAsync调用问题及非阻塞优化咨询

CUDA异步拷贝依赖参数的非阻塞解决方案

问题背景

在同一CUDA流中连续执行两次device-to-host方向的cudaMemcpyAsync调用:第一次拷贝一个int类型值(作为第二次拷贝的元素数量),第二次根据该值拷贝对应数量的double类型元素。核心代码如下:

cudaMemcpyAsync(h_size, d_size, sizeof(int), cudaMemcpyDeviceToHost, copy_stream);
cudaMemcpyAsync(h_array, d_array, h_size[0]*sizeof(double), 
                cudaMemcpyDeviceToHost, copy_stream);

实际运行发现,第二次拷贝使用的是第一次拷贝前h_size[0]的旧值。尽管同流操作会序列化执行,但cudaMemcpyAsync会在调用时立即确定所有参数值,而非流中实际执行时求值,导致第二次拷贝无法获取第一次拷贝更新后的h_size[0]。

使用cudaStreamSynchronize同步流可以解决问题,但会阻塞CPU线程,影响后续kernel的快速启动(代码包含循环执行kernel与通信的逻辑),因此需要非阻塞的解决方案。

非阻塞解决方案:CUDA流回调函数

利用cudaStreamAddCallback在第一次拷贝完成后,异步触发第二次拷贝。回调函数会在流中第一次拷贝执行完成后被调用,此时h_size[0]已经被更新,能正确计算第二次拷贝的字节数,且不会阻塞CPU线程。

修改后的communication函数代码

void communication(int * h_size, int const * d_size, 
                   double * h_array, double const * d_array,
                   cudaStream_t copy_stream,
                   cudaEvent_t compute_done,
                   cudaEvent_t copy_done)
{
    cudaStreamWaitEvent(copy_stream, compute_done, 0);
    // 第一次拷贝:将size从设备拷贝到主机
    cudaMemcpyAsync(h_size, d_size, sizeof(int), cudaMemcpyDeviceToHost, copy_stream);

    // 定义回调函数:在第一次拷贝完成后执行第二次拷贝
    auto copyCallback = [](cudaStream_t stream, cudaError_t status, void* data) {
        if (status != cudaSuccess) return;
        auto args = static_cast<struct CopyArgs*>(data);
        // 此时h_size已经更新,使用最新值计算拷贝大小
        size_t copy_bytes = args->h_size[0] * sizeof(double);
        cudaMemcpyAsync(args->h_array, args->d_array, copy_bytes, 
                        cudaMemcpyDeviceToHost, stream);
        cudaEventRecord(args->copy_done, stream);
        delete args;
    };

    // 封装回调需要的参数
    struct CopyArgs {
        int* h_size;
        const double* d_array;
        double* h_array;
        cudaEvent_t copy_done;
    };
    auto args = new CopyArgs{h_size, d_array, h_array, copy_done};

    // 将回调添加到流中,在第一次拷贝完成后执行
    cudaStreamAddCallback(copy_stream, copyCallback, args, 0);
}

方案说明

  1. 回调函数会在流中第一次cudaMemcpyAsync执行完成后被调用,此时h_size[0]已经是设备端更新后的值,能正确计算第二次拷贝的字节数。
  2. CPU提交完回调后即可继续执行后续代码(比如启动新的kernel),不会被阻塞,满足"尽快启动后续kernel"的需求。
  3. 注意在回调中处理参数的内存管理,避免内存泄漏(示例中用new分配参数结构体,回调内delete释放)。

原完整代码

#include <cuda_runtime.h>
#include <iostream>

__global__
void kernel(double* d_array, int* d_size, int const k)
{
    int i = blockIdx.x * blockDim.x + threadIdx.x;
    if (i==0) {
        printf("in kernel: d_size = %i\n", *d_size);
        if (k==0) {
            *d_size = 5;
        } else if (k==1) {
            *d_size = 3;
        } else {
            *d_size = 4;
        }
        printf("in kernel: d_size set to %i\n", *d_size);
        for (int i=0; i<3000; ++i) {
            __nanosleep(1000000U);
        }
    }
    d_array[i] += (1.0+(2.0+k)*i);
}

void communication(int * h_size, int const * d_size, 
                   double * h_array, double const * d_array,
                   cudaStream_t copy_stream,
                   cudaEvent_t compute_done,
                   cudaEvent_t copy_done)
{
    cudaStreamWaitEvent(copy_stream, compute_done, 0);
    // cudaMemcpyAsync in the next line will start as soon as the kernel execution is done
    cudaMemcpyAsync(h_size, d_size, sizeof(int), cudaMemcpyDeviceToHost, copy_stream);

    // Solution 1: 
    //cudaStreamSynchronize(copy_stream);
    // Problem: the next kernel won't start before this point

    // the next cudaMemcpyAsync is on the same stream as the previous one --> they serialize
    cudaMemcpyAsync(h_array, d_array, 
                    h_size[0]*sizeof(double), cudaMemcpyDeviceToHost, copy_stream);
    cudaEventRecord(copy_done, copy_stream);
}

void cpuWork(int const k, double const * h_array, int const * h_size)
{
    // here we check the values in h_array
    std::cout << "From kernel launch " << k << std::endl;
    for (int i=0; i<7; ++i) {
        std::cout << "h_array[" << i << "] = " << h_array[i] << std::endl;
    }
    std::cout << "h_size[" << k%2 << "] = " << h_size[k%2] << std::endl;
}

int main() {

    int bufferSize = 1<<18;
    int blockSize = 256;
    int numBlocks = (bufferSize + blockSize - 1)/blockSize;

    int *h_size;
    cudaMallocHost((void**)&h_size, 2*sizeof(int));
    h_size[0] = 2;
    h_size[1] = 7; 
    
    int *d_size;
    cudaMalloc((void**)&d_size, 2*sizeof(int));
    cudaMemcpy(d_size, h_size, 2*sizeof(int), cudaMemcpyHostToDevice);

    double *h_array;
    cudaMallocHost((void**)&h_array, bufferSize*sizeof(double));
    for (int i=0; i<bufferSize; ++i) {h_array[i] = 0.14;}

    double *d_array;
    cudaMalloc((void**)&d_array, 2*bufferSize*sizeof(double));
    cudaMemcpy(d_array, h_array, bufferSize*sizeof(double), cudaMemcpyHostToDevice);
    cudaMemcpy(&d_array[bufferSize], h_array, bufferSize*sizeof(double), cudaMemcpyHostToDevice);

    cudaStream_t compute_stream[2], copy_stream;
    cudaEvent_t compute_done[2], copy_done;

    cudaStreamCreateWithFlags(&compute_stream[0], cudaStreamNonBlocking);
    cudaStreamCreateWithFlags(&compute_stream[1], cudaStreamNonBlocking);
    cudaStreamCreateWithFlags(&copy_stream, cudaStreamNonBlocking);
    
    cudaEventCreateWithFlags(&compute_done[0], cudaEventDisableTiming);
    cudaEventCreateWithFlags(&compute_done[1], cudaEventDisableTiming);
    cudaEventCreateWithFlags(&copy_done, cudaEventDisableTiming);
    cudaDeviceSynchronize();
    
    // Launch the first kernel outside the FOR loop
    kernel<<<numBlocks, blockSize, 0, compute_stream[0]>>>(d_array, &d_size[0], 0);
    cudaEventRecord(compute_done[0], compute_stream[0]);
    
    int numCycles = 3;

    for (int k=1; k<numCycles; ++k) {
        
        // Communicate D-->H from the previous kernel launch, indexed k-1
        communication(&h_size[(k-1)%2], &d_size[(k-1)%2], 
                      h_array, &d_array[((k-1)%2)*bufferSize],
                      copy_stream, compute_done[(k-1)%2], copy_done);

        // Goal is to have this kernel start "asap" (even before the previous kernel has finished,
        // if the device has the capacity)
        kernel<<<numBlocks, blockSize, 0, compute_stream[k%2]>>>(
            &d_array[(k%2)*bufferSize], &d_size[k%2], k);
        cudaEventRecord(compute_done[k%2], compute_stream[k%2]);

        // Now do stuff on the CPU with h_array received from the (k-1)-th kernel launch
        cudaEventSynchronize(copy_done); // make sure communication is over
        cpuWork(k-1, h_array, h_size);
    }

    // Get results from the last cycle
    int k = numCycles-1;
    communication(&h_size[k%2], &d_size[k%2], 
                  h_array, &d_array[(k%2)*bufferSize],
                  copy_stream, compute_done[k%2], copy_done);

    // Now do work on the CPU with h_array received from the last kernel launch
    cudaEventSynchronize(copy_done); // make sure communication is over
    cpuWork(k, h_array, h_size);

    // clean
    for (int i=0; i<2; ++i) {cudaEventDestroy(compute_done[i]); 
                             cudaStreamDestroy(compute_stream[i]);}
    cudaEventDestroy(copy_done); cudaStreamDestroy(copy_stream);
    cudaFree(d_array); cudaFree(d_size);
    cudaFreeHost(h_array); cudaFreeHost(h_size);

    return 0;
}

预期与实际输出

预期结果:
From kernel launch 0
h_array[0] = 1.14 ok
h_array[1] = 3.14 ok
h_array[2] = 5.14 NO prints 0.14 means that h_size[0] was 2, not 5, at the 2nd cudaMemcpyAsync
h_array[3] = 7.14 NO prints 0.14 :
h_array[4] = 9.14 NO prints 0.14 :
h_array[5] = 0.14 ok
h_array[6] = 0.14 ok
h_size[0] = 5 ok
From kernel launch 1
h_array[0] = 1.14 ok
h_array[1] = 4.14 ok
h_array[2] = 7.14 ok
h_array[3] = 7.14 NO prints 10.14, means that h_size[1] was 7, not 3, at the 2nd cudaMemcpyAsync
h_array[4] = 9.14 NO prints 13.14,
h_array[5] = 0.14 NO prints 16.14,
h_array[6] = 0.14 NO prints 19.14,
h_size[1] = 3
From kernel launch 2
h_array[0] = 2.14 ok
h_array[1] = 8.14 ok
h_array[2] = 14.14 ok
h_array[3] = 20.14 ok
h_array[4] = 9.14 NO prints 26.14, showing that h_size[0] was 5, not 4, at the 2nd cudaMemcpyAsync
h_array[5] = 0.14 NO prints 16.14,
h_array[6] = 0.14 NO prints 19.14,
h_size[0] = 4

内容的提问来源于stack exchange,提问作者Samuel

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.30 11:29:49