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

CUDA流中第2轮迭代的数据传输与计算无法并发问题求助

CUDA流并发问题:首块计算与次块传输无法并行

问题描述

实现了一个基础CUDA程序,流程为:

  • 将数据加载到CPU页锁内存(Pinned)
  • 分块异步传输到GPU,每个数据块对应独立CUDA流
  • 对每个数据块执行计算

异常现象:第2个数据块的HostToDevice传输与第1个数据块的计算无法并发,二者顺序执行,但后续块的传输与对应计算均能正常并行。

已尝试方案:

  • 测试5次和10次迭代(有trace截图)
  • 使用nvcc 11.7和12.3版本编译
  • 改用3流模式(1个HostToDevice流、1个DeviceToHost流、1个计算流),问题依旧

测试环境

  • GPU:16GB V100
  • CUDA版本:11.7 / 12.3

原始代码

#include <bits/stdc++.h>
#include "driver_types.h"
#include "cuda.h"
#include "cuda_runtime.h"
#include "device_launch_parameters.h"
using namespace std;

#define checkCudaErrors(err)                                                     \
    do                                                                           \
    {                                                                            \
        if (err != cudaSuccess)                                                  \
        {                                                                        \
            std::cerr << "CUDA error at " << __FILE__ << ":" << __LINE__ << ": " \
                      << cudaGetErrorString(err) << std::endl;                   \
            exit(EXIT_FAILURE);                                                  \
        }                                                                        \
    } while (0)

__global__ void intialize_mark(char *mark, uint32_t num_e)
{
    uint32_t j;
    j = blockIdx.y * gridDim.x + blockIdx.x;
    j = j * 512 + threadIdx.x;
    if (j >= num_e)
        return;
    mark[j] = 0;
}

int main()
{
    cudaFree(nullptr);
    uint32_t numIterations = 10;
    std::ios_base::sync_with_stdio(false);
    std::cin.tie(0);
    std::cout.tie(0);

    uint32_t numNodes = 1000000;
    uint64_t numEdges = 20000000;

    uint64_t *edgelist;
    cudaHostAlloc(&edgelist, numEdges * sizeof(uint64_t), cudaHostAllocDefault);

    std::random_device rd;
    std::mt19937_64 gen(rd());
    std::uniform_int_distribution<uint64_t> dis;
    for (uint64_t i = 0; i < numEdges; i++)
    {
        edgelist[i] = dis(gen);
    }

    uint32_t num_threads = 512;                   // -> number of threads per block
    uint32_t num_blocks_n = (numNodes / 512) + 1; // -> number of blocks for nodes
    uint32_t num_blocks_e = (numEdges / 512) + 1; // -> number of blocks for edges
    uint32_t nny = (num_blocks_n / 1000) + 1;     // -> y dimension for nodes
    uint32_t nnx = 1000;                          // -> x dimension for nodes
    uint32_t ney = (num_blocks_e / 1000) + 1;     // -> y dimension for edges
    uint32_t nex = 1000;                          // -> x dimension for edges

    dim3 grid_n(nnx, nny);        // -> grid for nodes
    dim3 grid_e(nex, ney);        // -> grid for edges
    dim3 threads(num_threads, 1); // -> threads per block

    uint64_t numEdgesIteration = (numEdges + numIterations - 1) / numIterations; // -> number of edges per iteration

    char *d_mark;
    char *mask;
    uint64_t *d_edgeList1;
    uint64_t *d_edgeList2;

    checkCudaErrors(cudaMalloc(&d_mark, (numEdgesIteration) * sizeof(char)));
    checkCudaErrors(cudaMalloc(&mask, (numEdgesIteration) * sizeof(char)));

    checkCudaErrors(cudaMalloc(&d_edgeList1, (numEdgesIteration) * sizeof(uint64_t)));
    checkCudaErrors(cudaMalloc(&d_edgeList2, (numEdgesIteration) * sizeof(uint64_t)));

    cudaStream_t cudaStreamArr[numIterations];
    for (int i = 0; i < numIterations; i++)
    {
        cudaStreamCreate(&cudaStreamArr[i]);
    }
    uint64_t currentNumEdges;
    for (uint32_t i = 0; i <= numIterations; i++)
    {
        if (i < numIterations)
        {
            if ((min(numEdges, (i + 1) * numEdgesIteration) - (i)*numEdgesIteration) > 0)
            {
                checkCudaErrors(cudaMemcpyAsync(d_edgeList2,
                                                edgelist + (i)*numEdgesIteration,
                                                (min(numEdges, (i + 1) * numEdgesIteration) - (i)*numEdgesIteration) * sizeof(uint64_t),
                                                cudaMemcpyHostToDevice,
                                                cudaStreamArr[i]));
            }
        }

        if (i > 0)
        {
            currentNumEdges = min(numEdges, (i)*numEdgesIteration) - (i - 1) * numEdgesIteration;

            intialize_mark<<<grid_e, threads, 0, cudaStreamArr[i - 1]>>>(d_mark, currentNumEdges);

        }
        cudaDeviceSynchronize();
        swap(d_edgeList1, d_edgeList2);
    }
}

问题分析与修复方案

核心问题

代码中每次循环都调用cudaDeviceSynchronize(),该操作会阻塞主机端,等待所有CUDA流任务完成后才进入下一次迭代。这直接破坏了流的并发能力:

  • 第一次循环:仅启动第1块的传输,同步等待完成
  • 第二次循环:先启动第2块的传输,再启动第1块的计算,再次同步等待所有任务完成
  • 后续迭代的并发效果只是因为任务执行时间重叠,但本质上同步操作限制了并行潜力

此外还有两个潜在问题:

  1. 主机端交换设备指针swap(d_edgeList1, d_edgeList2)无同步保证,可能导致后续流操作访问错误指针
  2. 未销毁创建的CUDA流,存在资源泄漏

修复后代码

#include <bits/stdc++.h>
#include "driver_types.h"
#include "cuda.h"
#include "cuda_runtime.h"
#include "device_launch_parameters.h"
using namespace std;

#define checkCudaErrors(err)                                                     \
    do                                                                           \
    {                                                                            \
        if (err != cudaSuccess)                                                  \
        {                                                                        \
            std::cerr << "CUDA error at " << __FILE__ << ":" << __LINE__ << ": " \
                      << cudaGetErrorString(err) << std::endl;                   \
            exit(EXIT_FAILURE);                                                  \
        }                                                                        \
    } while (0)

__global__ void intialize_mark(char *mark, uint32_t num_e)
{
    uint32_t j;
    j = blockIdx.y * gridDim.x + blockIdx.x;
    j = j * 512 + threadIdx.x;
    if (j >= num_e)
        return;
    mark[j] = 0;
}

int main()
{
    cudaFree(nullptr);
    uint32_t numIterations = 10;
    std::ios_base::sync_with_stdio(false);
    std::cin.tie(0);
    std::cout.tie(0);

    uint32_t numNodes = 1000000;
    uint64_t numEdges = 20000000;

    uint64_t *edgelist;
    checkCudaErrors(cudaHostAlloc(&edgelist, numEdges * sizeof(uint64_t), cudaHostAllocDefault));

    std::random_device rd;
    std::mt19937_64 gen(rd());
    std::uniform_int_distribution<uint64_t> dis;
    for (uint64_t i = 0; i < numEdges; i++)
    {
        edgelist[i] = dis(gen);
    }

    uint32_t num_threads = 512;                   
    uint64_t numEdgesIteration = (numEdges + numIterations - 1) / numIterations; 

    // 根据块大小动态计算grid维度
    auto get_grid = [num_threads](uint64_t count) {
        uint32_t num_blocks = (count + num_threads - 1) / num_threads;
        uint32_t ny = (num_blocks / 1000) + 1;
        uint32_t nx = std::min(1000u, num_blocks);
        return dim3(nx, ny);
    };

    char *d_mark;
    uint64_t *d_edgeList[numIterations]; // 为每个块分配独立设备内存,避免指针冲突

    checkCudaErrors(cudaMalloc(&d_mark, numEdgesIteration * sizeof(char)));
    for (int i = 0; i < numIterations; i++) {
        checkCudaErrors(cudaMalloc(&d_edgeList[i], numEdgesIteration * sizeof(uint64_t)));
    }

    cudaStream_t cudaStreamArr[numIterations];
    for (int i = 0; i < numIterations; i++)
    {
        checkCudaErrors(cudaStreamCreate(&cudaStreamArr[i]));
    }

    // 批量启动所有块的异步传输
    for (uint32_t i = 0; i < numIterations; i++)
    {
        uint64_t start = i * numEdgesIteration;
        uint64_t end = std::min((i+1)*numEdgesIteration, numEdges);
        uint64_t size = end - start;
        if (size == 0) continue;

        checkCudaErrors(cudaMemcpyAsync(d_edgeList[i],
                                        edgelist + start,
                                        size * sizeof(uint64_t),
                                        cudaMemcpyHostToDevice,
                                        cudaStreamArr[i]));
    }

    // 为每个块启动计算,确保传输完成后执行
    for (uint32_t i = 0; i < numIterations; i++)
    {
        uint64_t size = std::min((i+1)*numEdgesIteration, numEdges) - i*numEdgesIteration;
        if (size == 0) continue;

        dim3 grid_e = get_grid(size);
        cudaEvent_t transfer_done;
        checkCudaErrors(cudaEventCreate(&transfer_done));
        checkCudaErrors(cudaEventRecord(transfer_done, cudaStreamArr[i]));
        checkCudaErrors(cudaStreamWaitEvent(cudaStreamArr[i], transfer_done, 0));

        intialize_mark<<<grid_e, num_threads, 0, cudaStreamArr[i]>>>(d_mark, size);
        checkCudaErrors(cudaGetLastError());
        checkCudaErrors(cudaEventDestroy(transfer_done));
    }

    // 等待所有任务完成后再清理资源
    checkCudaErrors(cudaDeviceSynchronize());

    // 释放所有资源
    checkCudaErrors(cudaFree(d_mark));
    for (int i = 0; i < numIterations; i++) {
        checkCudaErrors(cudaFree(d_edgeList[i]));
        checkCudaErrors(cudaStreamDestroy(cudaStreamArr[i]));
    }
    checkCudaErrors(cudaFreeHost(edgelist));

    return 0;
}

关键修复点

  1. 移除循环内的全局同步:仅在所有任务提交完成后执行一次全局同步
  2. 为每个块分配独立设备内存:避免指针交换导致的任务冲突,简化流依赖管理
  3. 添加流事件同步:确保每个块的计算任务在对应传输完成后才启动
  4. 动态计算grid维度:根据每个块的实际数据量生成合适的grid大小,避免资源浪费
  5. 完善资源清理:销毁CUDA流、释放设备/主机内存,避免泄漏

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 12:35:00