#include <stdio.h>
#include <cuda_runtime.h>
#define CUDA_CHECK(call) do { cudaError_t e = call; if(e != cudaSuccess) { \
fprintf(stderr, "CUDA Error: %s\n", cudaGetErrorString(e)); exit(1); } } while(0)
__global__ void compute_kernel(float* data, int n, float factor) {
int idx = threadIdx.x + blockIdx.x * blockDim.x;
if (idx < n) {
float val = data[idx];
for (int i = 0; i < 100; i++) {
val = val * factor + 1.0f;
}
data[idx] = val;
}
}
void serial_execution(float* h_data, float* d_data, int n, int num_chunks) {
int chunk_size = n / num_chunks;
int chunk_bytes = chunk_size * sizeof(float);
for (int i = 0; i < num_chunks; i++) {
CUDA_CHECK(cudaMemcpy(d_data + i * chunk_size,
h_data + i * chunk_size,
chunk_bytes, cudaMemcpyHostToDevice));
compute_kernel<<<chunk_size/256, 256>>>(d_data + i * chunk_size, chunk_size, 0.99f);
CUDA_CHECK(cudaMemcpy(h_data + i * chunk_size,
d_data + i * chunk_size,
chunk_bytes, cudaMemcpyDeviceToHost));
}
}
void stream_execution(float* h_data, float* d_data, int n, int num_streams) {
int stream_size = n / num_streams;
int stream_bytes = stream_size * sizeof(float);
cudaStream_t* streams = new cudaStream_t[num_streams];
for (int i = 0; i < num_streams; i++) {
cudaStreamCreate(&streams[i]);
}
float** h_chunk = new float*[num_streams];
for (int i = 0; i < num_streams; i++) {
cudaMallocHost(&h_chunk[i], stream_bytes);
for (int j = 0; j < stream_size; j++) {
h_chunk[i][j] = h_data[i * stream_size + j];
}
}
for (int i = 0; i < num_streams; i++) {
CUDA_CHECK(cudaMemcpyAsync(d_data + i * stream_size,
h_chunk[i],
stream_bytes,
cudaMemcpyHostToDevice,
streams[i]));
compute_kernel<<<>stream_size/256, 256, 0, streams[i]>>>(
d_data + i * stream_size, stream_size, 0.99f);
CUDA_CHECK(cudaMemcpyAsync(h_chunk[i],
d_data + i * stream_size,
stream_bytes,
cudaMemcpyDeviceToHost,
streams[i]));
}
for (int i = 0; i < num_streams; i++) {
CUDA_CHECK(cudaStreamSynchronize(streams[i]));
}
for (int i = 0; i < num_streams; i++) {
cudaStreamDestroy(streams[i]);
cudaFreeHost(h_chunk[i]);
}
delete[] streams;
delete[] h_chunk;
}
int main() {
const int N = 1 << 22;
const int num_chunks = 4;
printf("Stream并发Pipeline演示\n");
printf("数据大小: %.2f MB\n", N * sizeof(float) / 1024.0 / 1024.0);
float *h_data = (float*)malloc(N * sizeof(float));
float *d_data;
CUDA_CHECK(cudaMalloc(&d_data, N * sizeof(float)));
for (int i = 0; i < N; i++) h_data[i] = float(i);
cudaEvent_t start, stop;
cudaEventCreate(&start); cudaEventCreate(&stop);
float ms;
cudaEventRecord(start);
serial_execution(h_data, d_data, N, num_chunks);
cudaEventRecord(stop); cudaEventSynchronize(stop);
cudaEventElapsedTime(&ms, start, stop);
printf("串行执行: %.2f ms\n", ms);
cudaEventRecord(start);
stream_execution(h_data, d_data, N, num_chunks);
cudaEventRecord(stop); cudaEventSynchronize(stop);
cudaEventElapsedTime(&ms, start, stop);
printf("Stream并发: %.2f ms\n", ms);
cudaFree(d_data); free(h_data);
return 0;
}