六、多线程、并发与网络通信共 8 项一句话总览进程和线程是并发执行的两种基本载体进程间通信解决隔离环境下的数据交换线程同步解决共享地址空间下的安全协作Socket 把本地 IPC 扩展到网络TCP/UDP 决定传输语义高性能推理服务则通过事件循环、线程池、任务队列和流式返回把网络 I/O 与计算执行解耦。知识点关系图并发基础 ├─ 6. 进程 vs 线程 │ ├─ 1. 进程间通信 IPC │ ├─ pipe / FIFO │ ├─ message queue │ ├─ shared memory semaphore │ ├─ signal │ └─ socket │ ├─ 2. 多进程服务器 vs 多线程服务器 ├─ 3. Socket 与其他 IPC 的差异 ├─ 4. TCP / UDP / TCP 重传 │ ├─ 7. 共享内存安全与同步 ├─ 8. thread / mutex / condition_variable / atomic │ └─ 5. 高性能网络服务与流式输出 ├─ epoll / io_uring ├─ Reactor ├─ 推理队列 └─ SSE / WebSocket / chunked1. 进程间如何通信核心进程默认拥有独立地址空间不能直接访问彼此内存因此需要操作系统提供 IPCInter-Process Communication进程间通信机制。常见方式包括管道、命名管道、消息队列、共享内存、信号量、信号、Socket 和内存映射文件。1.1 常见 IPC 方式方式通信范围特点典型场景匿名管道pipe有亲缘关系进程半双工、字节流、生命周期随进程shell 管道命名管道FIFO同一主机任意进程文件系统中有路径本地生产者消费者消息队列同一主机按消息收发有边界本地任务消息共享内存同一主机最快需要同步大块数据交换信号量同一主机本身不传业务数据负责同步互斥、计数信号signal同一主机异步通知信息量小终止、重载、定时器Socket本机或跨网络双向、客户端/服务器模型网络服务、本地服务内存映射mmap同一主机文件或匿名内存映射文件共享、零拷贝1.2 匿名管道#include unistd.h #include iostream int main() { int fd[2]; pipe(fd); // fd[0] 读端fd[1] 写端 pid_t pid fork(); if (pid 0) { close(fd[0]); const char msg[] hello parent; write(fd[1], msg, sizeof(msg)); close(fd[1]); } else { close(fd[1]); char buf[64]{}; read(fd[0], buf, sizeof(buf)); std::cout buf \n; close(fd[0]); } }管道默认半双工数据像水流一样没有消息边界。1.3 命名管道 FIFOmkfifo /tmp/my_fifo一个进程写int fd open(/tmp/my_fifo, O_WRONLY); write(fd, task, 4);另一个进程读int fd open(/tmp/my_fifo, O_RDONLY); char buf[64]; read(fd, buf, sizeof(buf));1.4 共享内存共享内存把同一块物理内存映射到多个进程的虚拟地址空间进程 A 虚拟地址空间 进程 B 虚拟地址空间 ↓ ↓ └────── 同一块物理内存 ──────┘它避免了内核缓冲区和用户缓冲区之间的数据复制因此是最快的 IPC。但共享内存本身不提供同步A 写时 B 可能同时读需要信号量、互斥锁或环形队列配合。1.5 消息队列消息队列以消息为单位消息队列 ├─ type1: request A ├─ type2: request B └─ type1: request C它天然保留消息边界还可以按消息类型选择性读取但不适合超大块高频数据。1.6 信号信号是异步通知#include csignal #include iostream void handler(int sig) { std::cout receive signal sig \n; } int main() { signal(SIGTERM, handler); raise(SIGTERM); }信号只能携带很少信息处理函数中能安全调用的函数也有限不适合传输复杂数据。1.7 SocketSocket 既可用于本机AF_UNIX也可用于网络AF_INET/AF_INET6同一主机pipe/FIFO/共享内存/AF_UNIX socket 跨主机 只能使用网络 Socket选型可以按以下顺序只是父子进程传递文本 → pipe 本机任意进程、简单字节流 → FIFO 需要消息边界 → message queue 大块高频数据 → shared memory 同步原语 跨机器或未来可能网络化 → socket 异步事件通知 → signal/eventfd工程中最常见组合是“共享内存传大数据 信号量或 eventfd 通知 Socket 传控制命令”兼顾吞吐和易用性。2. 网络编程中多进程与多线程并发服务器有什么区别核心多进程服务器为每个连接或工作单元分配独立进程地址空间隔离、稳定性强但进程创建和切换成本高、共享数据困难多线程服务器在同一进程内并发处理连接线程共享内存和文件描述符通信轻量但需要锁、条件变量等同步机制一个线程触发严重错误可能拖垮整个进程。2.1 多进程模型典型流程主进程 listen/accept ├─ fork → 子进程 A处理连接 1 ├─ fork → 子进程 B处理连接 2 └─ fork → 子进程 C处理连接 3#include sys/socket.h #include unistd.h void run_server(int listen_fd) { while (true) { int conn accept(listen_fd, nullptr, nullptr); pid_t pid fork(); if (pid 0) { close(listen_fd); handle_connection(conn); _exit(0); } close(conn); } }优点进程地址空间隔离一个连接进程崩溃不直接破坏其他连接。不需要复杂锁因为默认不共享用户态内存。适合权限隔离、安全边界强的服务。可利用多 CPU进程调度由内核完成。缺点fork比线程创建重。进程上下文切换成本更高TLB、页表切换更明显。进程间共享状态需要 IPC。大量进程占用内存和文件描述符更多。2.2 多线程模型一个进程 ├─ 主线程accept ├─ 工作线程 1连接 A ├─ 工作线程 2连接 B └─ 工作线程 3连接 C 共享堆、全局变量、打开的文件描述符、代码段 私有栈、寄存器、线程局部存储、信号掩码#include thread #include vector void run_thread_server(int listen_fd) { std::vectorstd::thread workers; while (true) { int conn accept(listen_fd, nullptr, nullptr); workers.emplace_back([conn] { handle_connection(conn); }); // 实际工程不会无限创建线程通常使用线程池并 detach/join 管理 } }优点创建和切换成本低于进程。线程共享内存任务队列、缓存、连接池访问方便。线程间通信不需要内核 IPC。适合 I/O 等待多、需要共享连接池或模型实例的服务。缺点数据竞争、死锁、条件变量误唤醒等问题复杂。一个线程出现未捕获严重异常或非法内存访问整个进程可能退出。线程数过多会带来调度开销。一个阻塞操作可能占用工作线程影响整体吞吐。2.3 对比表维度多进程多线程地址空间独立共享创建成本高低切换成本高较低通信方式IPC共享变量、锁、队列崩溃隔离强弱调试难度隔离清晰并发问题复杂资源占用高较低共享缓存困难容易安全隔离更适合较弱2.4 现代服务器常用模型实际高性能服务通常不是“一连接一进程/线程”而是主 Reactoraccept 新连接 子 Reactorepoll 管理已连接 socket 线程池处理计算任务 推理/业务 Worker执行重计算多进程适合做故障域隔离例如多个 worker 进程各自管理模型进程内部再用多线程处理网络和批量计算。这样结合二者优点进程隔离风险线程降低通信成本。3. Socket 与其他通信方式有什么不同核心Socket 是通用通信端点既可以用于同一台机器也可以跨网络既支持可靠字节流也支持数据报。其他 IPC 方式通常只在本机有效并且在连接模型、数据边界、寻址能力和扩展性上各有限制。3.1 Socket 的基本特征Socket 是文件描述符但它比普通文件多了通信语义int fd socket(AF_INET, SOCK_STREAM, 0); bind(fd, ...); listen(fd, backlog); int conn accept(fd, ...); send(conn, ...); recv(conn, ...);客户端int fd socket(AF_INET, SOCK_STREAM, 0); connect(fd, ...); send(fd, ...);3.2 Socket 类型类型语义SOCK_STREAM面向连接、可靠、有序字节流通常对应 TCPSOCK_DGRAM无连接数据报通常对应 UDPSOCK_SEQPACKET面向连接、有消息边界、有序可靠SOCK_RAW原始报文可访问更底层协议AF_UNIX本机进程间 SocketAF_INETIPv4 网络 Socket3.3 与管道对比管道进程 A ──write── [pipe] ──read── 进程 B默认单向。主要用于父子或亲缘进程。不能跨机器。字节流无网络地址。Socket客户端 A ── IP Port ── 网络协议栈 ── 服务端 B可双向。可跨主机。有明确客户端/服务器模型。可选择 TCP/UDP 等协议。本机双向管道式通信也可以用socketpairint sv[2]; socketpair(AF_UNIX, SOCK_STREAM, 0, sv);3.4 与共享内存对比共享内存速度最快但它只解决“数据放哪里”不解决对方何时可读数据是否写完是否需要通知跨机器传输字节流可靠到达。Socket 会经过协议栈存在拷贝和协议开销但提供连接管理、路由、可靠传输、端口寻址和跨主机能力。3.5 与消息队列对比消息队列保留消息边界适合本地消息分发Socket 的 TCP 是字节流不保留应用层消息边界send 两次send(ab) send(cd) 接收端可能是 一次收到 abcd 或分两次收到 ab、cd因此 TCP 应用必须自己设计长度头、分隔符或协议解析状态机。UDP 则保留数据报边界一次recvfrom对应一个报文但不保证可靠到达。3.6 对比表通信方式跨主机双向消息边界连接模型典型用途pipe否默认单向无亲缘进程shell 管道FIFO否双向但需约定无文件路径本机服务消息队列否可收发有队列 key本地任务消息共享内存否双向由程序定义内存 key大块数据信号否单向通知无PID异步事件UNIX Socket本机双向取决于类型路径本机高性能服务网络 Socket是双向TCP 无/UDP 有IP端口网络通信3.7 选择原则只在本机、数据量极大 → 共享内存 轻量通知 本机简单控制通道 → AF_UNIX socket 需要跨网络 → TCP/UDP socket 需要广播/低延迟、可接受丢失 → UDP 需要可靠有序 → TCPSocket 的独特价值在于统一了本机和远程通信模型。一个服务先用AF_UNIX做本地部署后续需要远程访问时可以较平滑地切换到AF_INET这是管道和共享内存做不到的。4. TCP 和 UDP 的区别是什么TCP 什么时候会重传核心TCP 是面向连接、可靠、有序、面向字节流的协议UDP 是无连接、不保证可靠和顺序、面向数据报的协议。TCP 在未按时收到确认、收到重复确认或通过 SACK 发现报文丢失时会重传。4.1 TCP 与 UDP 对比维度TCPUDP是否连接需要三次握手建立连接无连接可靠性可靠到达、重传丢失不保证到达顺序保证按序交付不保证顺序数据形式字节流无消息边界数据报有边界流量/拥塞控制有无首部开销通常 20 字节起8 字节通信方式一对一为主可单播、广播、组播典型场景HTTP、RPC、文件传输DNS、语音、视频、游戏状态延迟较高但稳定低但可能丢包乱序4.2 TCP 连接建立与断开三次握手客户端 → SYN → 服务端 客户端 ← SYNACK ← 服务端 客户端 → ACK → 服务端 连接建立四次挥手主动方 → FIN 被动方 ← ACK 被动方 → FIN 主动方 ← ACK4.3 TCP 为什么要重传IP 层不保证报文一定到达网络可能丢包、乱序、重复。TCP 通过序号和确认号判断数据是否到达发送方发送1 2 3 4 5 接收方确认ACK 4 表示 4 之前都收到了 如果 4 一直没有 ACK发送方重传 44.4 超时重传 RTO发送每个报文段时启动重传定时器发送 segment ├─ RTO 时间内收到 ACK → 取消定时器 └─ RTO 到期仍未收到 ACK → 重传该 segmentRTO 不是固定值而是根据 RTT 动态估计SRTT平滑往返时间 RTTVARRTT 波动 RTO ≈ SRTT 4 × RTTVAR网络抖动越大RTO 越保守避免过早误重传但 RTO 太长又会拖慢恢复。4.5 快速重传如果接收方收到乱序报文会重复确认最后一个按序收到的序号发送1 2 3 4 5 3 丢失 接收 ACK 2 接收 ACK 2 接收 ACK 2 接收 ACK 2发送方收到 3 个重复 ACK通常不等 RTO立即重传缺失报文这叫快速重传3 个 duplicate ACK → Fast Retransmit → Fast Recovery4.6 SACK 选择性确认没有 SACK 时接收方只能告诉发送方“我连续收到哪里”。SACK 可以告诉对方哪些块已经收到已收到[1-100] [150-200] 缺失[101-149] 发送方只需重传缺失块现代 Linux 还使用 RACK、TLP 等机制基于发送时间和最新确认更精确地判断丢失减少虚假重传。4.7 UDP 不重传UDP 发送后不维护连接状态sendto(sock, data, len, 0, addr, sizeof(addr));丢了就丢了如果业务需要可靠性必须在应用层自己实现序号、ACK、重传和拥塞控制例如 QUIC 就是基于 UDP 实现可靠传输。4.8 选型文件、控制命令、网页请求、数据库协议 → TCP 实时语音、直播、在线游戏位置同步 → UDP 或应用层可靠 UDP 需要流式输出文本 → TCP 承载 HTTP/SSE/WebSocket总结TCP 的可靠性不是“网络不会丢”而是通过序号、确认、重传、流量控制和拥塞控制在不可靠 IP 层之上提供可靠字节流。重传主要由 RTO 超时、重复 ACK、SACK/RACK 丢包判断触发。5. C 高性能网络服务如何支撑大模型推理和流式输出核心把网络 I/O、请求排队、模型推理和结果回传拆成不同阶段网络线程使用 epoll/io_uring 等事件驱动机制处理大量连接推理线程通过有界队列接收任务模型逐 token 生成后通过 SSE、WebSocket 或 HTTP chunked 流式返回。关键是解耦、限流、背压、取消、批量调度和非阻塞写。5.1 总体架构客户端 │ HTTP/SSE/WebSocket ▼ 主 Reactoraccept 新连接 ▼ I/O 线程epoll 读写、协议解析 ▼ 有界请求队列mutex condition_variable ▼ 推理 Worker / 模型线程池 │ 逐 token 生成 ▼ 每连接输出队列 ▼ eventfd / epoll 通知 I/O 线程 ▼ 非阻塞 send 分片返回客户端5.2 为什么不能一个请求一个线程大模型推理耗时可能数秒到数十秒。如果每个连接创建一个线程长连接很多时线程数爆炸大量线程阻塞等待模型结果锁竞争和上下文切换严重无法统一做批量推理和优先级调度。因此应采用事件驱动 有限 worker。5.3 网络层epollint epfd epoll_create1(0); epoll_event ev{}; ev.events EPOLLIN | EPOLLET; ev.data.fd listen_fd; epoll_ctl(epfd, EPOLL_CTL_ADD, listen_fd, ev); while (true) { epoll_event events[1024]; int n epoll_wait(epfd, events, 1024, 100); for (int i 0; i n; i) { if (events[i].data.fd listen_fd) { accept_new_connection(); } else { handle_read_write(events[i].data.fd); } } }边缘触发 ET 通常要求 socket 设置为非阻塞并循环读到EAGAIN避免遗漏事件。5.4 流式输出协议SSE 响应头HTTP/1.1 200 OK Content-Type: text/event-stream Cache-Control: no-cache Connection: keep-alive data: {token:你} data: {token:好} data: [DONE]服务端每生成一个 token就推送一个 chunkvoid send_sse_token(int fd, std::string_view token) { std::string chunk data: std::string(token) \n\n; write_nonblock(fd, chunk); }WebSocket 更适合双向交互例如客户端随时停止生成HTTP chunked 适合普通流式响应。5.5 推理队列与背压队列必须有界#include condition_variable #include mutex #include queue template class Request class BoundedQueue { private: std::queueRequest q; std::mutex mtx; std::condition_variable not_full; std::condition_variable not_empty; std::size_t capacity; public: explicit BoundedQueue(std::size_t cap) : capacity(cap) {} bool push(Request req, std::chrono::milliseconds timeout) { std::unique_lock lock(mtx); if (!not_full.wait_for(lock, timeout, [] { return q.size() capacity; })) { return false; // 队列满触发限流或返回稍后重试 } q.push(std::move(req)); not_empty.notify_one(); return true; } Request pop() { std::unique_lock lock(mtx); not_empty.wait(lock, [] { return !q.empty(); }); Request req std::move(q.front()); q.pop(); not_full.notify_one(); return req; } };无界队列在流量突增时会耗尽内存因此必须设置最大排队长度、超时时间和拒绝策略。5.6 推理侧连续批处理请求 A生成 token 中 请求 B首 token 请求 C生成 token 中 ↓ Scheduler 把可运行请求组成一个 batch ↓ 一次模型前向产生多个请求的下一个 token这比为每个请求单独执行模型更能利用 GPU/NPU。5.7 输出回写生成 token 后不应由推理线程直接阻塞写 socket否则慢客户端会拖慢模型线程。更合理的是推理线程生成 token → 放入连接输出队列 → eventfd 通知 I/O 线程epoll 检测可写 → 尽量发送 → 剩余数据保留5.8 取消与超时用户断开连接时应让调度器把该请求标记为取消避免继续浪费算力struct RequestContext { std::atomicbool canceled{false}; std::queuestd::string output; };5.9 关键指标QPS 和并发连接数队列长度和排队等待时间首 token 延迟 TTFTtoken 间延迟 ITL吞吐 tokens/s取消率、超时率、拒绝率socket 写缓冲积压。总结高性能流式推理服务不是单纯“网络框架快”就够了而是网络事件循环、有界队列、批量调度、流式协议、背压取消和输出非阻塞回写共同配合的结果。6. 进程和线程的区别是什么核心进程是资源分配的基本单位拥有独立地址空间线程是 CPU 调度的基本单位存在于进程内部同一进程的线程共享代码、堆、全局变量和文件描述符但各自拥有栈和寄存器上下文。6.1 结构对比进程 ├─ 独立虚拟地址空间 ├─ 代码段 ├─ 数据段 ├─ 堆 ├─ 打开的文件 ├─ 权限和信号处理表 └─ 至少一个线程 进程内线程 ├─ 线程 1私有栈、寄存器、线程局部存储 ├─ 线程 2私有栈、寄存器、线程局部存储 └─ 共享代码段、堆、全局数据、文件描述符6.2 资源差异维度进程线程地址空间独立同进程内共享资源文件、内存、信号表独立共享进程资源栈每个进程有主栈每线程独立栈创建开销大小切换开销大较小通信IPC共享内存需同步崩溃影响通常只影响本进程可能使整个进程退出调度内核调度内核线程也由内核调度隔离性强弱6.3 进程示例#include unistd.h #include iostream int main() { pid_t pid fork(); if (pid 0) { std::cout child\n; } else if (pid 0) { std::cout parent\n; } }fork后子进程获得父进程资源的副本。现代系统通常使用写时复制 COW页面先共享真正修改时才复制。6.4 线程示例#include iostream #include thread void worker(int id) { std::cout thread id \n; } int main() { std::thread t1(worker, 1); std::thread t2(worker, 2); t1.join(); t2.join(); }两个线程可以直接访问同一全局变量#include atomic std::atomicint counter{0}; void add() { for (int i 0; i 1000; i) { counter.fetch_add(1); } }如果使用普通int就会产生数据竞争。6.5 上下文切换区别进程切换通常需要保存寄存器/程序计数器 切换页表基址 刷新或影响 TLB 切换内核栈、调度信息 恢复新进程上下文线程切换时同进程线程共享页表因此 TLB 和地址空间切换成本更低但仍要保存寄存器、栈指针和线程调度信息所以线程切换不是“零成本”。6.6 通信区别进程通信需要内核参与或共享映射pipe / message queue / socket / shared memory线程可以直接读写同一变量线程 A 写 queue 线程 B 读 queue 但必须用 mutex/cv/atomic 保证正确性6.7 故障隔离一个进程访问非法地址操作系统终止该进程多进程模型中其他进程仍可运行。一个线程访问非法地址会向整个进程发送致命信号所有线程一起结束。因此高可用服务常把强隔离任务拆成多个进程再在进程内使用线程池。6.8 选择原则需要安全隔离、独立权限、故障不扩散 → 多进程 需要共享缓存、低通信成本、大量 I/O 并发 → 多线程 CPU 密集计算 → 进程/线程数接近硬件并行度 I/O 密集等待 → 事件驱动 少量线程 高可用服务 → 多进程 worker 进程内线程池总结进程回答“资源归谁、如何隔离”线程回答“代码从哪里并发执行、如何调度”。二者不是替代关系现代服务器通常组合使用。7. 共享内存安全吗有什么措施保证核心共享内存本身并不安全它只是让多个进程映射同一块物理内存。如果多个执行流同时读写且没有同步就会出现数据竞争、读到半完成数据、指令重排可见性问题甚至未定义行为。保证安全的关键是同步、内存序、单写者原则、生命周期管理和只共享可跨进程安全使用的数据。7.1 为什么不安全时刻 T1进程 A 正在写 header.length 10 时刻 T2进程 B 同时读取 header.length 和 data 结果B 可能读到旧长度 新数据或新长度 旧数据即使硬件是多核缓存一致也不代表程序自动线程/进程安全。缓存一致保证缓存行最终同步但 C 内存模型要求程序使用原子、锁或其他同步原语建立 happens-before 关系。7.2 典型风险数据竞争多个进程同时写同一字段。撕裂读写字段不是原子写入读端读到中间状态。顺序错误写入数据后再写 ready 标志但没有内存屏障读端可能先看到 ready。队列头尾竞争生产者和消费者更新索引不同步。生命周期问题一个进程卸载共享内存另一个仍访问。错误共享伪共享高频写变量落在同一 cache line导致性能下降。跨进程 C 对象问题含虚函数、指针、std::string、STL 容器的对象不能简单放到共享内存中跨进程使用。7.3 措施一互斥锁 条件变量进程共享互斥锁需要使用进程共享属性#include pthread.h #include sys/mman.h struct SharedData { pthread_mutex_t mtx; int value; }; void init(SharedData* data) { pthread_mutexattr_t attr; pthread_mutexattr_init(attr); pthread_mutexattr_setpshared(attr, PTHREAD_PROCESS_SHARED); pthread_mutexattr_setrobust(attr, PTHREAD_MUTEX_ROBUST); pthread_mutex_init(data-mtx, attr); } void write_value(SharedData* data, int value) { pthread_mutex_lock(data-mtx); >7.4 措施二原子变量对于计数器、标志位、队列索引可使用原子操作#include atomic #include cstdint struct RingBuffer { std::atomicuint32_t head; std::atomicuint32_t tail; char data[4096]; };跨进程使用std::atomic时应确认它在目标平台是 lock-freestatic_assert(std::atomicuint32_t::is_always_lock_free);不同 CPU 架构、编译器和标准库下非 lock-free 的 atomic 可能包含进程内锁对象不适合直接放入共享内存。7.5 措施三SPSC 环形队列单生产者单消费者环形队列可以做到无锁生产者只写 write_index 消费者只写 read_index 双方都读取对方索引 data[write] 写完并发布后再推进 write_index简化代码template class T, std::size_t N class SPSCRing { private: T buffer[N]; std::atomicstd::size_t write{0}; std::atomicstd::size_t read{0}; public: bool push(const T value) { std::size_t w write.load(std::memory_order_relaxed); std::size_t r read.load(std::memory_order_acquire); std::size_t next (w 1) % N; if (next r) return false; // 满 buffer[w] value; write.store(next, std::memory_order_release); return true; } bool pop(T value) { std::size_t r read.load(std::memory_order_relaxed); std::size_t w write.load(std::memory_order_acquire); if (r w) return false; // 空 value buffer[r]; read.store((r 1) % N, std::memory_order_release); return true; } };7.6 措施四内存序和屏障写数据 → release 发布索引 读端 acquire 看到索引 → 保证能看到此前写入的数据不要随意使用memory_order_relaxed传递数据依赖。没有把握时使用默认seq_cst正确优先再根据性能测试优化。7.7 措施五数据设计尽量只放平凡可复制类型、定长数组、整数索引。不放裸指针因为不同进程映射地址可能不同如需内部引用使用共享内存内偏移量。不放虚对象因为虚指针指向各进程自己的代码段布局未必一致。对 cache line 对齐避免 false sharingstruct alignas(64) CacheLineCounter { std::atomicuint64_t value{0}; };7.8 安全措施总结风险措施多写竞争mutex、semaphore、单写者状态通知condition variable、semaphore、eventfd索引/标志atomic acquire/release高频队列SPSC/MPSC 环形缓冲持锁崩溃robust mutex 或无锁设计伪共享cache line 对齐指针失效使用偏移量而非绝对指针生命周期引用计数、统一创建/销毁进程结论共享内存“快”但不会自动“对”。它适合传输大块数据控制信息和读写顺序必须由同步机制保证。8. C 线程、mutex、condition_variable、atomic 如何用于推理队列核心std::thread提供工作线程std::mutex保护共享任务队列std::condition_variable在线程没有任务时睡眠、有任务时唤醒std::atomic管理无锁状态标志例如停止信号和计数器。四者组合可以实现一个典型的生产者-消费者推理队列。8.1 队列结构网络线程 / 调用线程 │ submit(request) ▼ std::queueRequest ← mutex 保护 │ notify_one ▼ 条件变量唤醒等待的 Worker │ ├─ Worker 1pop → infer → callback ├─ Worker 2pop → infer → callback └─ Worker 3pop → infer → callback8.2 完整示例#include atomic #include condition_variable #include functional #include iostream #include mutex #include queue #include thread #include vector struct InferRequest { int id; std::vectorfloat input; }; struct InferResponse { int id; std::vectorfloat output; }; class InferenceQueue { private: std::queueInferRequest jobs; std::mutex mtx; std::condition_variable cv; std::vectorstd::jthread workers; std::atomicbool stopping{false}; std::atomicint active_count{0}; using Callback std::functionvoid(InferResponse); Callback callback; InferResponse model_infer(const InferRequest req) { // 真实场景中调用模型 runtime InferResponse resp; resp.id req.id; resp.output req.input; return resp; } void worker_loop() { while (true) { InferRequest req; { std::unique_lockstd::mutex lock(mtx); cv.wait(lock, [] { return stopping.load() || !jobs.empty(); }); // 停止且队列已空时退出 if (stopping.load() jobs.empty()) { return; } req std::move(jobs.front()); jobs.pop(); active_count.fetch_add(1); } InferResponse resp model_infer(req); if (callback) { callback(std::move(resp)); } active_count.fetch_sub(1); } } public: explicit InferenceQueue(int worker_num, Callback cb) : callback(std::move(cb)) { for (int i 0; i worker_num; i) { workers.emplace_back([this] { worker_loop(); }); } } ~InferenceQueue() { shutdown(); } void submit(InferRequest req) { if (stopping.load()) { throw std::runtime_error(queue is stopping); } { std::lock_guardstd::mutex lock(mtx); jobs.push(std::move(req)); } cv.notify_one(); } void shutdown() { bool expected false; if (!stopping.compare_exchange_strong(expected, true)) { return; } cv.notify_all();; // std::jthread 析构会自动 join } std::size_t pending_size() { std::lock_guardstd::mutex lock(mtx); return jobs.size(); } };使用int main() { InferenceQueue queue(4, [](InferResponse resp) { std::cout finish request resp.id \n; }); for (int i 0; i 10; i) { queue.submit(InferRequest{i, {1.0f, 2.0f}}); } }8.3 为什么 wait 必须带谓词条件变量可能发生虚假唤醒没有 notifywait 也可能返回。惊群效应多个线程被唤醒但只有一个任务。任务被其他线程取走醒来后队列仍可能为空。因此必须写成cv.wait(lock, [] { return stopping || !jobs.empty(); });它等价于循环检查while (!stopping jobs.empty()) { cv.wait(lock); }8.4 锁的范围不要在持锁状态下执行模型推理// 错误推理期间一直占锁其他线程无法取任务 lock.lock(); auto req jobs.front(); jobs.pop(); auto result model.infer(req); lock.unlock();正确方式是只在操作队列时持锁{ lock; 取出任务; } 释放锁后执行耗时推理;8.5 notify_one 与 notify_allnotify_one唤醒一个 worker适合每个任务只需要一个消费者。notify_all唤醒所有等待线程适合 shutdown、多个条件同时变化或多个任务批量可用。8.6 atomic 的使用边界适合 atomic 的内容std::atomicbool stop{false}; std::atomicint queue_size{0}; std::atomicuint64_t request_id{0};不适合把整个复杂请求变成 atomic。复杂对象仍应放入队列并由 mutex 保护。8.7 有界队列与限流生产环境不能使用无限std::queue应增加最大长度bool try_submit(InferRequest req, std::chrono::milliseconds timeout) { std::unique_lock lock(mtx); if (!cv_not_full.wait_for(lock, timeout, [] { return jobs.size() max_size; })) { return false; } jobs.push(std::move(req)); cv_not_empty.notify_one(); return true; }8.8 推理队列最佳实践1. 网络线程只解析和入队不做重计算 2. Worker 数量按模型并行度和硬件资源设置 3. 队列必须有界防止流量冲垮内存 4. 锁只保护队列不保护模型执行 5. 停止标志用 atomic退出用 notify_all 6. 请求 ID、计数器用 atomic 7. 输出结果通过回调、future 或输出队列返回 8. 流式生成时每个 token 可以推送到连接输出队列总结thread负责并发执行mutex保证共享队列互斥访问condition_variable避免空转轮询atomic轻量管理状态。它们组合起来就是高性能服务中最常见的任务队列骨架。本部分总结题号主题核心结论1IPC管道、消息队列、共享内存、Socket 各有适用范围共享内存最快但需同步2多进程/多线程进程重隔离线程重共享和低开销现代服务常组合使用3Socket唯一天然支持跨网络的通用 IPCTCP 是字节流UDP 是数据报4TCP/UDPTCP 可靠有序UDP 低延迟无保证RTO、重复 ACK、SACK 触发重传5流式服务epoll 事件循环 有界队列 worker SSE/WebSocket 流式返回6进程/线程进程是资源分配单位线程是调度单位7共享内存本身不安全需要锁、原子、内存序、环形队列和生命周期管理8推理队列thread mutex condition_variable atomic 构成生产者消费者模型整体记忆主线隔离进程 并行线程 通信IPC / Socket 可靠TCP 序号确认与重传 高性能事件驱动 线程池 有界队列 正确性mutex / atomic / condition_variable / 内存序 流式输出token 分片、非阻塞写、背压与取消
