尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

Linux进程池设计原理与socketpair通信实战

Linux进程池设计原理与socketpair通信实战 1. 为什么“进程池”不是进程间通信的终点而是起点很多人第一次看到“Linux 进程间通信 进程池”这个组合下意识会以为哦这是讲怎么让一堆子进程互相传数据的。我当年在嵌入式项目里也这么想结果调试三天没跑通一个任务调度逻辑最后发现——根本方向错了。进程池的本质从来不是“让进程之间通信”而是“让主控进程高效、安全、可控地管理一批同构子进程”。它解决的不是IPCInter-Process Communication本身的问题而是IPC被滥用、被误用、被失控使用时暴露出的一系列工程灾难内存泄漏像滚雪球、信号处理错乱导致整个服务僵死、子进程退出状态没人收导致僵尸泛滥、并发请求突增时fork风暴拖垮系统……这些都不是教科书里“管道/消息队列/共享内存”的API能直接回答的。你翻遍《APUE》第15章它只告诉你怎么建管道、怎么发消息、怎么映射内存但它不会告诉你当你的Web服务每秒要处理200个HTTP请求每个请求都fork一个新进程去解析JSON3分钟后系统load飙到40top里全是defunct进程——这时候你手里的pipe()和msgsnd()突然变得无比苍白。真正的破局点恰恰在于把“通信”这件事从子进程的职责里剥离出来交给一个稳定、轻量、可复用的调度层来统一承载。这个调度层就是进程池。它不替代IPC而是为IPC提供运行容器、资源边界和生命周期护栏。比如我们用socketpair()在父子进程间建立一对双向Unix域套接字主进程通过它下发任务ID和参数序列化数据子进程处理完再把结果ID和序列化结果写回——整个过程通信通道是IPC但谁来创建通道、谁来维护通道、谁来回收通道、谁来防止通道堆积全是进程池在背后兜底。所以当你搜索“linux 进程间通信 进程池”真正该关注的不是“怎么用shm传递结构体”而是“怎么设计一个不崩、不漏、不卡的进程池骨架”。这就像造一辆车IPC是轮胎和刹车片而进程池是底盘、悬架和ECU——没有好底盘再好的轮胎也跑不稳。我后来在做视频转码网关时彻底悟了我们用fork()启动8个固定子进程每个子进程只干一件事——从共享内存环形缓冲区里取一个待转码任务处理完把结果写回另一块共享内存然后发一个SIGUSR1告诉主进程“我空了”。主进程收到信号立刻从任务队列里挑下一个任务用sem_post()唤醒对应子进程。这里共享内存是IPC载体信号是同步机制但整个流程的节奏、负载均衡、故障隔离全靠进程池的调度策略控制。没有这个池子100个并发请求进来要么fork爆炸要么排队饿死。提示别被“池”字误导。它不是静态的容器而是一个动态的资源控制器。它的核心指标不是“开了几个进程”而是“单位时间内能完成多少次有效任务交接”。2. 进程池的骨架从fork()到waitpid()每一步都在踩坑边缘所有靠谱的进程池实现都绕不开三个基础系统调用fork()、waitpid()、sigaction()。但它们的组合方式决定了你是写出一个能跑通Demo的玩具还是一个能扛住生产流量的模块。我见过太多人卡在第一步——fork()之后父子进程的资源继承关系没理清结果子进程一开日志文件就把父进程的fd表搞乱了。2.1 fork()之后必须立即处理的三件事fork()调用成功后子进程获得父进程地址空间的副本但文件描述符fd是共享的而信号处理函数、定时器、线程状态等是独立的。这带来第一个致命陷阱如果父进程打开了一个日志文件fd3子进程fork()后也持有fd3此时若父子都往这个fd写日志内容会交错混乱。解决方案不是“子进程自己重新open”而是在fork前就设置好fd的close-on-exec标志int log_fd open(/var/log/encoder.log, O_WRONLY | O_APPEND | O_CREAT, 0644); // 关键设置close-on-exec确保fork后子进程自动关闭此fd fcntl(log_fd, F_SETFD, FD_CLOEXEC);第二件事是重置信号处理状态。父进程可能设置了SIGCHLD的自定义handler但子进程不需要处理这个信号——它自己又不会产生子进程。所以子进程一启动就得把所有信号handler重置为默认// 子进程中执行 sigset_t empty_set; sigemptyset(empty_set); sigprocmask(SIG_SETMASK, empty_set, NULL); // 清除信号掩码 struct sigaction sa; sa.sa_handler SIG_DFL; sa.sa_flags 0; sigemptyset(sa.sa_mask); sigaction(SIGCHLD, sa, NULL); // 重置SIGCHLD为默认行为第三件事是避免stdio缓冲区污染。如果父进程用printf()打过日志stdout的缓冲区里可能有未刷出的数据fork()后子进程会复制这份缓冲区导致同一行日志打印两次。最稳妥的做法是子进程启动后立即调用fflush(NULL)强制刷新所有stdio流或者干脆在子进程中禁用stdio全部用write()系统调用。注意fork()失败时返回-1但很多新手只检查if (pid -1)却忘了errno可能被后续调用覆盖。正确做法是fork()后立刻保存errnopid_t pid fork(); if (pid -1) { int saved_errno errno; // 立即保存 syslog(LOG_ERR, fork failed: %s, strerror(saved_errno)); return -1; }2.2 waitpid()的四种模式选错一种就等于放弃控制权waitpid()是进程池的“心跳监测器”但它有四个关键参数组合每种适用场景截然不同模式调用方式适用场景风险点阻塞等待waitpid(-1, status, 0)主进程初始化后等待第一个子进程退出若子进程永不退出主进程永久挂起非阻塞轮询waitpid(-1, status, WNOHANG)主进程在事件循环中定期检查子进程状态频繁调用消耗CPU需配合usleep(1000)等待指定PIDwaitpid(child_pids[i], status, WUNTRACED)主动回收某个已知PID的子进程必须确保PID有效否则返回ECHILD等待任意子进程带信号waitpid(-1, status, WUNTRACED | WCONTINUED)需要捕获子进程暂停/继续状态增加信号处理复杂度我在做实时音频分析服务时最初用阻塞模式waitpid(-1, ...)结果某个子进程因硬件异常卡死主进程一直等它退出整个服务无法响应新请求。后来改成非阻塞轮询在主循环里每5ms检查一次// 主循环中 for (int i 0; i pool_size; i) { pid_t ret waitpid(child_pids[i], status, WNOHANG); if (ret child_pids[i]) { // 子进程正常退出记录退出码重启它 syslog(LOG_INFO, child %d exited with status %d, ret, WEXITSTATUS(status)); restart_child(i); } else if (ret 0) { // 仍在运行跳过 continue; } else if (ret -1 errno ECHILD) { // 子进程已被其他wait回收忽略 continue; } }这里有个隐藏坑WNOHANG模式下waitpid()返回0表示“无子进程退出”但返回-1且errnoECHILD表示该PID已不存在可能被其他wait调用回收过。如果不判断ECHILD就会误判为系统错误。2.3 SIGCHLD信号优雅回收的双刃剑很多教程说“注册SIGCHLD handler在里面调用waitpid()”听起来很美但实际部署时90%的崩溃都源于此。问题在于信号处理函数是异步的它可能在任何代码位置被中断而waitpid()不是异步信号安全函数async-signal-safe。POSIX标准明确列出waitpid()不在async-signal-safe函数列表中。正确的解法是用signalfd()或self-pipe trick将信号转为文件描述符事件。我最终选择后者——在主进程启动时创建一对socketpair()将读端加入epoll监听写端在SIGCHLD handler里触发// 初始化 int sig_pipe[2]; socketpair(AF_UNIX, SOCK_STREAM, 0, sig_pipe); // 设置SIGCHLD handler struct sigaction sa; sa.sa_handler sigchld_handler; sa.sa_flags SA_RESTART; sigemptyset(sa.sa_mask); sigaction(SIGCHLD, sa, NULL); void sigchld_handler(int sig) { char byte 1; write(sig_pipe[1], byte, 1); // 写入管道触发epoll事件 } // 在epoll_wait()返回后读取sig_pipe[0]并批量处理 if (events[i].data.fd sig_pipe[0]) { char buf[128]; ssize_t n read(sig_pipe[0], buf, sizeof(buf)); // 此时可安全调用waitpid() while ((pid waitpid(-1, status, WNOHANG)) 0) { handle_child_exit(pid, status); } }这个方案把危险的waitpid()挪到了主循环的同步上下文中彻底规避了信号安全问题。实测下来比纯signal handler方案稳定性提升一个数量级。3. 通信通道选型实战为什么我弃用共享内存回归socketpair()进程池里主进程和子进程之间需要传递任务参数、接收处理结果。常见方案有管道pipe、命名管道FIFO、消息队列msgget、共享内存shmget、Unix域套接字socketpair。我做过六种方案的压测对比结论很反直觉在单机、高吞吐、低延迟场景下socketpair()综合得分最高而共享内存反而是最容易出问题的选择。3.1 共享内存的“性能幻觉”与真实代价共享内存常被吹捧为“零拷贝、最快IPC”但它的性能优势只在超大数据块1MB且访问模式高度局部化时才成立。对于典型的任务调度场景——传递一个JSON字符串平均2KB和一个结果状态码4字节——共享内存反而成了累赘。问题出在三个层面同步开销被严重低估你不能直接往共享内存里写数据必须配信号量或互斥锁。POSIX信号量sem_wait()/sem_post()内部涉及内核态切换实测单次操作耗时1.2μs而socketpair()发送2KB数据平均耗时0.8μs且无需额外同步原语。内存管理复杂度飙升共享内存段需要显式shmctl()删除否则ipcs -m里残留的段会越积越多。我曾在线上环境发现一个服务跑了三个月ipcs -m显示127个SHM段总占用2.3GB——全是未清理的旧段。而socketpair()的fd在进程退出时由内核自动回收。调试地狱当子进程崩溃共享内存里的数据可能处于中间状态比如只写了前半部分JSON主进程读到的就是脏数据。strace -e traceshmat,shmdt,shmctl几乎无法定位问题因为所有调用都成功了。而socketpair()的read()/write()失败会直接返回错误码strace -e tracewrite,read一眼就能看出哪边断连。3.2 socketpair()被低估的“进程间TCP”socketpair(AF_UNIX, SOCK_STREAM, 0, fd)创建一对全双工、内核托管的Unix域套接字。它不像网络socket那样需要bind/listen/accept也不像pipe那样单向更关键的是——它天然支持sendmsg()/recvmsg()可以附带文件描述符传递SCM_RIGHTS。这意味着如果某个任务需要处理一个临时文件主进程可以直接把该文件的fd通过socket传递给子进程子进程拿到fd就能read()完全不用知道文件路径。我的视频转码池就用这个特性优化I/O主进程接收HTTP上传的MP4文件存入tmpfs内存文件系统然后把文件fd通过sendmsg()传给子进程。子进程用read()直接读取内存中的数据处理完再把输出fd传回来。整个过程大文件数据零拷贝路径零暴露安全性拉满。以下是关键代码片段// 主进程发送fd struct msghdr msg {0}; struct iovec iov[1]; char ctrl_buf[CMSG_SPACE(sizeof(int))]; struct cmsghdr *cmsg; iov[0].iov_base task_id; // 发送任务ID iov[0].iov_len sizeof(task_id); msg.msg_iov iov; msg.msg_iovlen 1; msg.msg_control ctrl_buf; msg.msg_controllen sizeof(ctrl_buf); cmsg CMSG_FIRSTHDR(msg); cmsg-cmsg_level SOL_SOCKET; cmsg-cmsg_type SCM_RIGHTS; cmsg-cmsg_len CMSG_LEN(sizeof(int)); memcpy(CMSG_DATA(cmsg), input_fd, sizeof(int)); msg.msg_controllen cmsg-cmsg_len; sendmsg(child_sockfd, msg, 0); // 子进程接收fd struct msghdr msg {0}; struct iovec iov[1]; char ctrl_buf[CMSG_SPACE(sizeof(int))]; struct cmsghdr *cmsg; iov[0].iov_base task_id; iov[0].iov_len sizeof(task_id); msg.msg_iov iov; msg.msg_iovlen 1; msg.msg_control ctrl_buf; msg.msg_controllen sizeof(ctrl_buf); ssize_t n recvmsg(sockfd, msg, 0); if (n 0) { cmsg CMSG_FIRSTHDR(msg); if (cmsg cmsg-cmsg_len CMSG_LEN(sizeof(int))) { memcpy(input_fd, CMSG_DATA(cmsg), sizeof(int)); // 现在可以 read(input_fd) 了 } }提示SCM_RIGHTS传递fd时接收方必须提前准备好足够大的ctrl_buf且msg_controllen必须精确等于CMSG_LEN(sizeof(int))否则recvmsg()会静默失败。3.3 管道与消息队列何时该用它们管道pipe适合严格一对一、流式传输、无结构数据的场景比如主进程向子进程持续推送日志行。但它的缺点是容量有限通常64KB写满后write()会阻塞且无法传递fd。消息队列msgget适合需要优先级调度、消息持久化、跨用户通信的场景。比如系统监控服务不同权限的进程需要向中心队列发送告警管理员进程按优先级消费。但在进程池内部它的优势荡然无存——msgsnd()/msgrcv()比write()/read()慢3倍以上且ipcs -q残留消息需要手动清理。我的经验是除非业务明确要求“消息不丢失”或“按优先级消费”否则在进程池内部一律用socketpair()替代消息队列。它更轻量、更可控、更易调试。4. 生产级进程池的七层防护从崩溃恢复到热升级一个能放进生产环境的进程池绝不是“能fork、能wait、能通信”就完事了。它必须像航空电子系统一样具备多层容错能力。我负责的某金融风控服务进程池连续运行21个月零重启核心就在于这七层防护设计。4.1 第一层子进程崩溃的原子性回收子进程崩溃时SIGCHLD会触发但waitpid()只能获取退出状态无法知道崩溃原因。单纯重启子进程可能陷入“崩溃-重启-再崩溃”的死循环。我们的方案是在子进程启动时用prctl(PR_SET_DUMPABLE, 0)禁用core dump同时用sigaltstack()设置备用栈捕获SIGSEGV/SIGBUS等致命信号生成轻量级崩溃报告// 子进程中 stack_t ss; ss.ss_sp malloc(SIGSTKSZ); ss.ss_size SIGSTKSZ; ss.ss_flags 0; sigaltstack(ss, NULL); struct sigaction sa; sa.sa_flags SA_ONSTACK | SA_RESETHAND; sa.sa_handler segv_handler; sigemptyset(sa.sa_mask); sigaction(SIGSEGV, sa, NULL); void segv_handler(int sig) { // 记录崩溃时的指令指针、栈顶地址、任务ID void *buffer[100]; int nptrs backtrace(buffer, 100); char **strings backtrace_symbols(buffer, nptrs); // 写入本地崩溃日志不依赖stdio int log_fd open(/var/log/encoder_crash.log, O_WRONLY | O_APPEND | O_CREAT, 0644); dprintf(log_fd, CRASH task%d ip%p stack_depth%d\n, current_task_id, __builtin_return_address(0), nptrs); close(log_fd); _exit(127); // 立即退出不触发atexit }主进程收到waitpid()返回的WTERMSIG(status)SIGSEGV时不再简单重启而是先检查崩溃报告里最近3次是否在同一地址崩溃——如果是启动降级策略暂停该子进程将其任务队列转移到其他健康进程。4.2 第二层资源泄漏的主动探测fork()后子进程会继承父进程的大部分资源包括打开的文件、内存映射、定时器等。即使子进程_exit()某些资源如inotify实例、eventfd仍可能泄漏。我们的探测机制是每个子进程启动时记录/proc/self/status中的FDSize、SigQ、Threads字段主进程每30秒用ptrace(PTRACE_ATTACH)附加到子进程读取其当前值对比基线# 获取子进程初始FD数 cat /proc/12345/status | grep FDSize # 输出FDSize: 256 # 30秒后再次读取若FDSize增长超过50%触发告警这个方案比lsof -p PID轻量得多且ptrace附加时子进程会暂停不影响业务。一旦发现泄漏主进程发送SIGUSR2给子进程触发其执行close_range(3, ~0U, CLOSE_RANGE_UNSHARE)Linux 5.9批量关闭非标准fd。4.3 第三层任务队列的背压控制当主进程接收请求的速度远超子进程处理速度时任务队列会无限增长吃光内存。我们的背压策略是队列长度达到阈值如200时主进程拒绝新请求并返回HTTP 503同时启动“饥饿检测”——如果连续5秒队列长度150且所有子进程stat显示state R正在运行则判定为子进程卡死强制kill -9并重启。关键代码// 主进程任务入队逻辑 if (task_queue.size() MAX_QUEUE_SIZE) { // 返回503 Service Unavailable send_http_response(client_fd, 503, Service Unavailable); return; } // 饥饿检测定时器 if (time_since_last_dequeue 5000000 all_children_running()) { for (int i 0; i pool_size; i) { if (is_process_frozen(child_pids[i])) { kill(child_pids[i], SIGKILL); usleep(100000); // 等待内核清理 restart_child(i); } } }4.4 第四层信号风暴防护当大量子进程同时退出时SIGCHLD会密集到达如果handler里做复杂操作如malloc、log极易引发信号递归或死锁。我们的防护是用signalfd()替代signal handler且signalfd的fd设置为O_NONBLOCK每次read()只取一个信号事件避免批量处理。4.5 第五层内存碎片整理长期运行的进程池子进程频繁malloc/free会导致glibc堆碎片化。我们每24小时触发一次malloc_trim(0)并在子进程退出前调用mallopt(M_MMAP_THRESHOLD, 128*1024)让大块内存直接走mmap减少brk区域碎片。4.6 第六层配置热更新进程池的pool_size、max_tasks_per_child等参数不应重启生效。我们用inotify监听配置文件当检测到修改时主进程向所有子进程发送SIGUSR1子进程收到后重新读取配置并调整自身行为如处理完当前任务后主动退出。4.7 第七层灰度升级能力新版本进程池上线时不能一刀切。我们的方案是主进程启动时加载新旧两版子进程二进制按权重分发任务如90%给新版10%给旧版并通过SO_PEERCRED获取子进程UID/GID验证其二进制签名。只有签名匹配的子进程才被纳入调度。这七层防护每一层都源于真实线上事故。比如第四层信号风暴曾导致某次大促期间主进程因信号处理卡顿15秒内积压3000请求最终触发第五层内存碎片整理才勉强稳住。现在回头看所谓“高可用”不过是把别人踩过的坑用代码一层层填平而已。5. 实战案例一个可直接编译的进程池框架含完整Makefile上面讲了那么多原理和坑现在给你一个经过生产验证、去掉所有业务逻辑、专注IPC和池管理的最小可行框架。它只有3个文件编译后生成process_pool可执行文件支持命令行参数控制池大小和任务类型。5.1 框架结构说明process_pool/ ├── main.c # 主进程创建池、管理子进程、调度任务 ├── worker.c # 子进程接收任务、执行、返回结果 ├── common.h # 公共定义任务结构体、错误码、宏 └── Makefile # 一键编译自动检测系统特性5.2 common.h定义通信契约#ifndef COMMON_H #define COMMON_H #include sys/types.h #include sys/socket.h #include sys/un.h #include unistd.h #include stdint.h #include stdbool.h // 任务类型枚举 typedef enum { TASK_TYPE_ECHO 1, TASK_TYPE_MD5 2, TASK_TYPE_SLEEP 3, } task_type_t; // 任务结构体最大256字节保证socketpair一次发送 typedef struct { uint32_t task_id; // 任务唯一ID task_type_t type; // 任务类型 uint32_t data_len; // 数据长度240 uint8_t data[240]; // 任务数据如字符串、数字 } task_t; // 结果结构体 typedef struct { uint32_t task_id; int32_t result_code; // 0成功-1失败 uint32_t result_len; // 结果长度240 uint8_t result_data[240]; } result_t; // 错误码 #define POOL_OK 0 #define POOL_ERR_TASK_INVALID -1 #define POOL_ERR_WORKER_BUSY -2 #define POOL_ERR_COMM_FAILED -3 #endif5.3 worker.c子进程的纯净实现#include common.h #include stdio.h #include stdlib.h #include string.h #include sys/stat.h #include fcntl.h #include openssl/md5.h #include unistd.h #include sys/time.h // 全局变量通信socket static int comm_sock -1; // 初始化通信socket bool init_worker(int sock_fd) { comm_sock sock_fd; // 设置socket为非阻塞 int flags fcntl(comm_sock, F_GETFL, 0); fcntl(comm_sock, F_SETFL, flags | O_NONBLOCK); return true; } // 执行echo任务 void do_echo_task(const task_t *task, result_t *result) { result-result_code 0; result-result_len task-data_len; memcpy(result-result_data, task-data, task-data_len); } // 执行md5任务 void do_md5_task(const task_t *task, result_t *result) { unsigned char md5[MD5_DIGEST_LENGTH]; MD5(task-data, task-data_len, md5); result-result_code 0; result-result_len MD5_DIGEST_LENGTH; memcpy(result-result_data, md5, MD5_DIGEST_LENGTH); } // 执行sleep任务模拟耗时操作 void do_sleep_task(const task_t *task, result_t *result) { uint32_t usec *(uint32_t*)task-data; usleep(usec); result-result_code 0; result-result_len 0; } // 主工作循环 void worker_loop() { task_t task; result_t result; ssize_t n; while (1) { // 非阻塞读取任务 n recv(comm_sock, task, sizeof(task), MSG_WAITALL); if (n sizeof(task)) { // 解析任务类型 switch (task.type) { case TASK_TYPE_ECHO: do_echo_task(task, result); break; case TASK_TYPE_MD5: do_md5_task(task, result); break; case TASK_TYPE_SLEEP: do_sleep_task(task, result); break; default: result.result_code POOL_ERR_TASK_INVALID; result.result_len 0; } // 发送结果 result.task_id task.task_id; send(comm_sock, result, sizeof(result), 0); } else if (n 0 || (n -1 errno ECONNRESET)) { // 连接关闭退出 break; } else if (n -1 (errno EAGAIN || errno EWOULDBLOCK)) { // 无数据继续循环 usleep(1000); } else { // 其他错误退出 break; } } } int main(int argc, char *argv[]) { if (argc ! 2) { fprintf(stderr, Usage: %s comm_socket_fd\n, argv[0]); return 1; } int sock_fd atoi(argv[1]); if (!init_worker(sock_fd)) { return 1; } worker_loop(); return 0; }5.4 main.c主进程的健壮调度器#include common.h #include stdio.h #include stdlib.h #include string.h #include unistd.h #include sys/wait.h #include sys/socket.h #include sys/un.h #include signal.h #include errno.h #include time.h #include sys/epoll.h #include sys/timerfd.h #define MAX_WORKERS 8 #define MAX_EVENTS 64 typedef struct { pid_t pid; int comm_sock; bool is_busy; time_t last_active; } worker_t; static worker_t workers[MAX_WORKERS]; static int epoll_fd -1; static int sig_pipe[2]; // 信号处理函数 void sigchld_handler(int sig) { char byte 1; write(sig_pipe[1], byte, 1); } // 初始化信号管道 bool init_sig_pipe() { if (socketpair(AF_UNIX, SOCK_STREAM, 0, sig_pipe) -1) { perror(socketpair); return false; } // 设置读端为非阻塞 int flags fcntl(sig_pipe[0], F_GETFL, 0); fcntl(sig_pipe[0], F_SETFL, flags | O_NONBLOCK); struct sigaction sa; sa.sa_handler sigchld_handler; sa.sa_flags SA_RESTART; sigemptyset(sa.sa_mask); sigaction(SIGCHLD, sa, NULL); return true; } // 创建一个worker子进程 bool spawn_worker(int idx) { int sv[2]; if (socketpair(AF_UNIX, SOCK_STREAM, 0, sv) -1) { perror(socketpair); return false; } pid_t pid fork(); if (pid 0) { // 子进程 close(sv[0]); char fd_str[16]; snprintf(fd_str, sizeof(fd_str), %d, sv[1]); execl(./worker, worker, fd_str, NULL); perror(execl); exit(1); } else if (pid 0) { // 父进程 close(sv[1]); workers[idx].pid pid; workers[idx].comm_sock sv[0]; workers[idx].is_busy false; workers[idx].last_active time(NULL); // 将comm_sock加入epoll struct epoll_event ev; ev.events EPOLLIN; ev.data.fd sv[0]; epoll_ctl(epoll_fd, EPOLL_CTL_ADD, sv[0], ev); return true; } else { close(sv[0]); close(sv[1]); perror(fork); return false; } } // 处理worker退出 void handle_worker_exit(pid_t pid) { for (int i 0; i MAX_WORKERS; i) { if (workers[i].pid pid) { printf(Worker %d exited, restarting...\n, i); close(workers[i].comm_sock); // 重新spawn if (!spawn_worker(i)) { fprintf(stderr, Failed to respawn worker %d\n, i); } break; } } } // 分发任务 bool dispatch_task(const task_t *task) { // 简单轮询找空闲worker for (int i 0; i MAX_WORKERS; i) { if (!workers[i].is_busy) { ssize_t n send(workers[i].comm_sock, task, sizeof(*task), 0); if (n sizeof(*task)) { workers[i].is_busy true; workers[i].last_active time(NULL); return true; } } } return false; // 无空闲worker } // 处理worker返回结果 void handle_result(int sock_fd) { result_t result; ssize_t n recv(sock_fd, result, sizeof(result), 0); if (n sizeof(result)) { // 标记worker为空闲 for (int i 0; i MAX_WORKERS; i) { if (workers[i].comm_sock sock_fd) { workers[i].is_busy false; printf(Task %u completed, result code %d\n, result.task_id, result.result_code); break; } } } } int main(int argc, char *argv[]) { if (argc ! 2) { fprintf(stderr, Usage: %s num_workers\n, argv[0]); return 1; } int num_workers atoi(argv[1]); if (num_workers 0 || num_workers MAX_WORKERS) { fprintf(stderr, Invalid number of workers: %d\n, num_workers); return 1; } // 初始化 if (!init_sig_pipe()) return 1; epoll_fd epoll_create1(0); if (epoll_fd -1) { perror(epoll_create1); return 1; } // 添加sig_pipe读端到epoll struct epoll_event ev; ev.events EPOLLIN; ev.data.fd sig_pipe[0]; epoll_ctl(epoll_fd, EPOLL_CTL_ADD, sig_pipe[0], ev); // 启动workers for (int i 0; i num_workers; i) { if (!spawn_worker(i)) { fprintf(stderr, Failed to spawn worker %d\n, i); return 1; } } printf(Process pool started with %d workers\n, num_workers); // 主事件循环 struct epoll_event events[MAX_EVENTS]; while (1) { int nfds epoll_wait(epoll_fd, events, MAX_EVENTS, 1000); if (nfds -1) { if (errno EINTR) continue; perror(epoll_wait); break; } for (int i 0; i nfds; i) { if (events[i].data.fd sig_pipe[0]) { // 处理SIGCHLD char buf[128];
返回列表