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

跨进程线程乒乓通信无法扩展的原因探究

线程/进程乒乓通信的扩展性性能差异分析

本文探究线程/进程乒乓通信的扩展性问题:基于相同数据结构,设置两种测试场景:

  • 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 07:47:02