跨进程线程乒乓通信无法扩展的原因探究
线程/进程乒乓通信的扩展性性能差异分析
本文探究线程/进程乒乓通信的扩展性问题:基于相同数据结构,设置两种测试场景:
- CASE-1:两个进程各包含N个线程执行乒乓通信;
- CASE-2:启动2N个进程执行乒乓通信。
测试结果:
当仅使用2个进程时,两者吞吐量均为1.2 mops;但当线程/进程数增至2N时,CASE-1(两进程共96线程)的吞吐量仅为5mops,而CASE-2(192个进程)的吞吐量可达60mops。现分析该性能差异的原因,同时排查是否存在测试误差。
核心代码逻辑
ping-side 代码
local variable count = 0; while (true) { ping++; count++; while (count > 0) { if (!pong) sched_yield(); else break; } if (pong) { pong--; count--; } }
pong-side 代码
while (true) { while (!ping) { sched_yield(); } pong++; ping--; }
CASE-1 实现逻辑
pid_t t = fork(); if (t > 0) { for (int t = 0; t < N; t++) { thread([&](int tid){ ping-side operating on ping[tid]/pong[tid] }, t); } } else if (t == 0) { for (int t = 0; t < N; t++) { thread([&](int tid){ pong-side operating on ping[tid]/pong[tid] }, t); } }
CASE-2 实现逻辑
for (int t = 0; t < 2 * N; t++) { pid_t tid = fork(); if (tid == 0) { if (t < N)){ ping-side operating on ping[t % N]/pong[t % N] } else { pong-side operating on ping[t % N]/pong[t % N] } break; } else if (t == 0) { wait_all(); } }
测试环境与编译命令
测试环境为96核平台,编译运行命令:
g++ -O3 -o test test.cpp -lnuma -lpthread -lrt CASE-1 ./test 96 1000000 1 96 CASE-2 ./test 96 1000000 0 96
完整测试代码
#include <atomic> #include <iostream> #include <vector> #include <thread> #include <mutex> #include <sched.h> #include <numa.h> #include <unistd.h> #include <fcntl.h> #include <sys/wait.h> #include <sys/mman.h> #include <sys/stat.h> #define USE_THREAD 0 #define preempt() sched_yield() //usleep(1) //sched_yield class atomicPQ { alignas(64) std::atomic<uint64_t> ping{0}; alignas(64) std::atomic<uint64_t> pong{0}; public: atomicPQ() { ping.store(0); pong.store(0); } uint64_t ping_degrade(int num = -1) { return ping.fetch_add(num); } uint64_t pong_degrade(int num = -1) { return pong.fetch_add(num); } int64_t ping_fetch() { return ping.load(); } int64_t pong_fetch() { return pong.load(); } void ping_upgrade() { ping.fetch_add(1); } void pong_upgrade() { pong.fetch_add(1); } }; void ping(int id, int MAX_CORES, bool prod = true) { cpu_set_t mask; CPU_ZERO(&mask); CPU_SET(id % MAX_CORES, &mask); sched_setaffinity(0, sizeof(cpu_set_t), &mask); } void threadSwitch(int argc, char** argv) { int work_num = std::atoi(argv[1]); size_t total_round = std::atol(argv[2]); int type = std::atoi(argv[3]); int MAX_CORES = std::atoi(argv[4]); constexpr uint64_t max_queue_depth = 1; int commem = shm_open("fusecomplete", O_RDWR | O_CREAT | O_APPEND, S_IRUSR | S_IWUSR); for (int i = 0; i <= 512; i++) { char buf[4096]; memset(buf, 0, 4096); write(commem, buf, 4096); } std::vector<std::thread> producers, consumers, calculators; std::atomic<uint64_t> terminated{0}, total_count{0}, total_sum{0}, total_sumcnt{0}; auto start_time = std::chrono::high_resolution_clock::now(); pid_t retid = 1; atomicPQ* sem = (atomicPQ*) mmap(NULL, 1048576 * 2, PROT_READ | PROT_WRITE, MAP_SHARED, commem, 0); memset(sem, 0, 1048576 * 2); retid = fork(); if (retid == 0) { std::cout << "producer: " << getpid() << std::endl; for (int i = 0; i < work_num; i++) { producers.emplace_back([&](int tid) { ping(tid, MAX_CORES); size_t count = 0; int active = 0; for (int j = 0; j < total_round; j++) { sem[tid].ping_upgrade(); active++; while (active == max_queue_depth) { if (terminated.load() > 0) break; if (sem[tid].pong_fetch() == 0) { preempt(); } else break; } if (sem[tid].pong_fetch() > 0) { sem[tid].pong_degrade(); active--; } if (terminated > 0) break; count++; } while (active != 0) { // rest requests to be handled completely. if (sem[tid].pong_fetch() > 0) { sem[tid].pong_degrade(); active--; } else preempt(); } sem[work_num].ping_upgrade(); // communicate with consumer for quit notification terminated.fetch_add(1); total_count.fetch_add(count); }, i); } for (auto& p: producers) { p.join(); } for (auto& t: calculators) t.join(); auto end_time = std::chrono::high_resolution_clock::now(); uint64_t microseconds = std::chrono::duration_cast<std::chrono::microseconds>(end_time - start_time).count(); std::cout << work_num << " " << type << " " << total_round << " " << (double) total_count / microseconds << " " << total_count << " " << sizeof(atomicPQ) << " " << microseconds << " " << (double) total_sum.load() << " " << (double) total_sumcnt.load() << " " << (double) total_sumcnt.load() / microseconds << std::endl; } else if (retid > 0) { std::cout << "consumer: " << getpid() << std::endl; for (int i = 0; i < work_num; i++) { consumers.emplace_back([&](int tid) { ping(tid, MAX_CORES); int count = 0; for (int j = 0; j < total_round; j++) { while (sem[tid].ping_fetch() == 0) { if (sem[work_num].ping_fetch() == work_num) break; else { preempt(); } } if (sem[work_num].ping_fetch() == work_num) break; // must be consistent with line 208 to quit sem[tid].ping_degrade(); sem[tid].pong_upgrade(); count++; } while (sem[tid].ping_fetch() > 0) { sem[tid].ping_degrade(); sem[tid].pong_upgrade(); } }, i); } int status; pid_t wpid; while ((wpid = wait(&status)) > 0); for (auto& c: consumers) c.detach(); } shm_unlink("fusecomplete"); } void processSwitch(int argc, char** argv) { uint64_t max_cores = std::atoi(argv[1]); uint64_t total_count = std::atol(argv[2]); int MAX_CORES = std::atoi(argv[4]); const size_t parfactor = 2; constexpr uint64_t max_queue_depth = 1; std::mutex outlock; int commem = shm_open("fusecomplete", O_RDWR | O_CREAT | O_APPEND, S_IRUSR | S_IWUSR); for (int i = 0; i <= 512; i++) { char buf[4096]; memset(buf, 0, 4096); write(commem, buf, 4096); } atomicPQ* sem = (atomicPQ*) mmap(NULL, 1048576, PROT_READ | PROT_WRITE, MAP_SHARED, commem, 0); memset(sem, 0, 1048576); auto start_time = std::chrono::high_resolution_clock::now(); pid_t pid; int work_num = max_cores * parfactor / 2; for (int i = 0; i < work_num * 2; i++) { pid = fork(); if (pid == 0) { int64_t qidx = i % work_num, count = 0; if (i >= work_num) { // consumers ping(i, MAX_CORES); for (int k = 0; k < total_count; k++) { while (sem[qidx].ping_fetch() == 0) { if (sem[work_num].ping_fetch() == work_num) break; else { preempt(); } } if (sem[work_num].ping_fetch() == work_num) break; sem[qidx].ping_degrade(); sem[qidx].pong_upgrade(); count++; } while (sem[qidx].ping_fetch() > 0) { sem[qidx].ping_degrade(); sem[qidx].pong_upgrade(); } } else { // producers ping(i, MAX_CORES); int64_t active = 0; for (int k = 0; k < total_count; k++) { sem[qidx].ping_upgrade(); active++; while (active == max_queue_depth) { if (sem[work_num].ping_fetch() > 0) break; if (sem[qidx].pong_fetch() == 0) { preempt(); } else break; } if (sem[qidx].pong_fetch() > 0) { sem[qidx].pong_degrade(); active--; count++; } if (sem[work_num].ping_fetch() > 0) break; } while (active != 0) { // rest requests to be handled completely. if (sem[qidx].pong_fetch() > 0) { sem[qidx].pong_degrade(); active--; count++; } else { preempt(); } } sem[work_num].ping_upgrade(); // communicate with consumer for quit notification sem[work_num].pong_degrade(count); struct timespec tp; sched_rr_get_interval(gettid(), &tp); double tm = tp.tv_sec * 1000.0f + tp.tv_nsec / 1000000.0f; outlock.lock(); std::cout << i << " " << count << " " << sem[qidx].ping_fetch() << " " << sem[qidx].pong_fetch() << " " << sem[work_num * 2 + 1 + qidx].ping_fetch() << " " << sem[work_num * 2 + 1 + qidx].pong_fetch() << " " << tm << " " << gettid() << std::endl; outlock.unlock(); } break; } } if (pid > 0) { int status; pid_t wpid; while ((wpid = wait(&status)) > 0); auto end_time = std::chrono::high_resolution_clock::now(); struct timespec tp; sched_rr_get_interval(getpid(), &tp); double tm = tp.tv_sec * 1000.0f + tp.tv_nsec / 1000000.0f; uint64_t microseconds = std::chrono::duration_cast<std::chrono::microseconds>(end_time - start_time).count(); std::cout << "Throughput: " << (double) sem[max_cores * parfactor / 2].pong_fetch() / microseconds << " " << sem[max_cores * parfactor / 2].pong_fetch() << " " << tm << std::endl; shm_unlink("fusecomplete"); } } int main(int argc, char** argv) { int type = std::atoi(argv[3]); if (type == 0) processSwitch(argc, argv); else threadSwitch(argc, argv); }
性能差异原因分析
1. 调度与上下文切换开销
- CASE-1中,两个进程各有48个线程(共96线程),同一进程内的线程共享地址空间,内核调度时需处理线程间上下文切换,高并发下调度竞争加剧,
sched_yield()调用频率上升,大量CPU时间浪费在调度环节。 - CASE-2中,192个进程对应独立调度单元,96核平台可让每个核心绑定1-2个进程,进程执行上下文更稳定,
sched_yield()无效调用更少,调度开销被显著降低。
2. 缓存与内存一致性开销
- CASE-1中,同一进程内的多线程共享L1/L2缓存,频繁访问共享内存的原子变量会引发缓存行颠簸;即便用
alignas(64)保证缓存行独占,线程跨核心调度时,MESI缓存一致性协议的开销也会被放大。 - CASE-2中,进程绑定固定核心后,每个核心上的进程访问对应缓存行的频率稳定,缓存命中率更高,内存一致性开销被分散到不同核心,无进程内多线程的缓存竞争。
3. 同步与全局变量竞争
- CASE-1中,
terminated原子变量被所有生产者线程频繁读取,带来额外缓存一致性开销;进程内的线程join()操作也会增加同步成本。 - CASE-2中,每个进程独立,终止信号通过共享内存的特定变量传递,竞争范围极小,同步开销可忽略。
4. 测试时序差异
CASE-1中主进程需等待子进程退出后才detach消费者线程,导致消费者线程执行延迟;CASE-2中所有进程同时启动,执行时序更紧凑,无额外等待损耗。
测试误差排查
从代码与命令来看,测试逻辑基本一致,但需确认两点:
- 核心绑定:CASE-1中生产者与消费者线程的
tid范围均为0-95,会导致同一核心绑定多个线程;CASE-2中进程i范围0-191,每个核心绑定2个进程,核心利用率更充分。 - 统计逻辑:CASE-1通过
total_count原子变量统计,CASE-2通过共享内存变量统计,两者统计逻辑准确,无明显误差。
内容的提问来源于stack exchange,提问作者Li Wenhai
相关产品推荐
相关产品推荐

