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

Linux进程队列通信异常:master加sleep后user进程阻塞

问题描述

在Linux环境下实现了master主进程与user子进程的通信架构:master创建N个user进程,user通过System V消息队列向master请求预算。master创建子进程后进入无限循环等待所有子进程终止,无sleep时系统运行正常(仅存在数学逻辑待修复)。但为满足每秒打印状态的需求,在master循环中加入sleep(1)后出现异常:

  • user进程仅打印getBudget: [PID]和setBudget: [PID]日志,队列接收结果及后续操作日志无输出,仿佛setBudget函数未执行完毕;
  • 若仅保留user启动日志,该日志仅在进程终止时才打印,所有user进程陷入无限循环。

怀疑sleep影响了waitpid或msgrcv的执行,但按此推测user应能打印错误日志,实际却无相关输出。


master.c

userList = initSharedMemory(SH_KEY_USER,SH_USERS_LIST_SIZE);
nodeList = initSharedMemory(SH_KEY_NODE,SH_NODES_LIST_SIZE);

int queueID = msgget(KEY_QUEUE, IPC_CREAT | 0600);

char *argsUser[2] = {USER_NAME,NULL};
char *argsNode[2] = {NODE_NAME,NULL};

for (int i = 0; i < SO_NODES_NUM; i++)
{
    int pid = fork();
    switch (pid){
        case 0:
            if(execvp(NODE_NAME, argsNode)==-1){
                    perror("error in execvp: ");
                }
            exit(EXIT_FAILURE);
            break;
        case -1:
            printf("Error in main:  fork failed");
            break;
        default:
            printf("\nmaster: launch node %d (pid %d)\n",i,pid);
            nodeList[i] = pid;
            break;
    }
}
for (int i = 0; i < SO_USERS_NUM; i++)
    {
        int pid = fork();
        switch (pid){
            case 0:
                if(execvp(USER_NAME, argsUser)==-1){
                    perror("error in execvp: ");
                }
                exit(EXIT_FAILURE);
                break;
            case -1:
                printf("Error in main:  fork failed");
                break;
            default:
                printf("\nmaster: launch user %d (pid %d)\n",i,pid);
                userList[i] = pid;
                break;
        }
    }

struct msgUser mt_msg;
pid_t pidnow;
int num_bytes;


while (1)
{
    //resumePrint();
    sleep(1);
    //printf("\nciclo post sleep  %d   %d",user_done,node_done);
    if (user_done+node_done == SO_USERS_NUM+SO_NODES_NUM) {
        break;
    }

      if (user_done == SO_USERS_NUM)
    {
        for (int i = 0; i < SO_NODES_NUM; i++)
        {
            kill(nodeList[i], 5);
            sleep(2);
            kill(nodeList[i], 9);
            sleep(1);
        }
    }
    
    num_bytes = msgrcv(queueID, &mt_msg, MSGUSER_SIZE, MSG_BUDGET, IPC_NOWAIT);
    //printf("\nciclo post sleep  numbytes:%d",num_bytes);
    if (num_bytes >= 0)
    {
        receivedBudgetRequest(mt_msg,queueID);
    }

    pidnow= waitpid(-1, NULL, WNOHANG);
    //printf("\nciclo post sleep  pid_finish:%d ",pidnow);
    if(pidnow ==-1)perror("error in waitpid: ");
    if(pidnow > 0){
        for (int i = 0; i < SO_USERS_NUM; i++)
        {
            if (pidnow == userList[i])
            {
                
                printf("\nTerminato user: done:%d pid:%d\n", user_done, pidnow);
                ++user_done;
                break;
            }
        }
        for (int j = 0; j < SO_NODES_NUM; j++)
        {
            if (pidnow == nodeList[j])
            {
                
                printf("\nTerminato node: done:%d pid:%d\n", node_done, pidnow);
                ++node_done;
                break;
            }
        }
    }
  
}

user.c

int main(int argc, char const *argv[]){

nodeList = initSharedMemory(SH_KEY_NODE,SH_NODES_LIST_SIZE);

if (!initTransactionPool())
{
    printf("\nError in initTransaction\n");
    exit(EXIT_FAILURE);
}
int queueID = msgget(KEY_QUEUE, IPC_CREAT | 0600);
setQueueId(queueID);

attach_signals();
srand(getpid());
int finish = 0;

printf("\n sono l'user: %d",getpid()); // if i didn't put other print this will be pritned only at the end and not when the process start

while (finish < SO_RETRY)
{
    checkTransactionCompleted();
    getBudget();
    int result = setBudget();
    
    if(result == 0){
        printf("\n result setBudget: %d   user: %d",result,getpid()); // doesn't print 
        int val = makeTransaction();
        if(val == -1){
            ++finish;
        } 
    }
    int delay = rand() % SO_MAX_TRANS_GEN_NSEC / 1000000000 + SO_MIN_TRANS_GEN_NSEC / 1000000000;
    sleep(delay);
}
printf("\nuser terminato");
freeTransactionPool();
return 0;
}

库函数

void getBudget()
{
    printf("\ngetBudget: %d",getpid());
    struct msgUser messageUser;
    messageUser.mtype = MSG_BUDGET;
    messageUser.pid = getpid();
    if(msgsnd(queueID, &messageUser, MSGUSER_SIZE, 0)==-1){
        perror("\nerror in msgsnd getBudget");
    }
    
}
int setBudget()
{
    
    struct msgUser mt_msg;
    int num_bytes;
    totalBudget=SO_BUDGET_INIT;
    
    num_bytes = msgrcv(queueID, &mt_msg, MSGUSER_SIZE, id, 0);
    printf("\n setBudget: %d",getpid()); // it print this and anything else
    if(num_bytes == -1)perror("msgrcv error in setBudgt:");
    if (num_bytes > 0)
    {
        
        if((mt_msg.numero>SO_BUDGET_INIT*SO_USERS_NUM) || mt_msg.numero<-(SO_BUDGET_INIT*SO_USERS_NUM)) return mt_msg.numero;

        totalBudget += mt_msg.numero; 
        printf("\n received by master: %d   %d",totalBudget,getpid());
        
    }
    checkTransactionCompleted(queueID); // doesn't need if the master don't send anything
    int sum = sumMoneyTransactions();
    if((totalBudget +sum)<-998) {
            printf("\nError addition: %d %d",(totalBudget+sum),getppid());
            return -1;
    }
    totalBudget += sum;
    printf("\nBudget: %d   %d\n\n",totalBudget,getpid());
    return 0;
}

void receivedBudgetRequest(struct msgUser mt_msg,int queueID)
{
  int initialBudget = calculateBudgetFor(mt_msg.pid);
  struct msgUser messageMaster;
  messageMaster.mtype = mt_msg.pid; 
  messageMaster.pid = getpid();
  messageMaster.numero = initialBudget; 

  msgsnd(queueID, &messageMaster, MSGUSER_SIZE, 0);
  printf("\n mando messaggio ad user");  // doen't print
}
int sumMoneyTransactions()
{
    int tot = 0;
    for (int i = 0; i < sizePool; i++)
    {
        tot -= transaction_pool[i].money + transaction_pool[i].reward;
    }
    return tot;
}

问题分析与解决方案

核心问题1:master消息处理被延迟

master循环中先执行sleep(1),再调用msgrcv(IPC_NOWAIT),意味着user发送的预算请求最多要等待1秒才会被处理。而user在setBudget中是阻塞式调用msgrcv(最后一个参数为0),导致user进程卡在msgrcv处,后续代码完全无法执行——这就是为什么user只打印到setBudget: [PID]的原因。

核心问题2:stdout缓冲导致日志延迟

user的启动日志仅在进程终止时打印,是因为标准输出默认是行缓冲,该日志没有换行符,内容被留在缓冲区中,直到进程退出时才被统一刷新。

修复步骤

  1. 调整master循环执行顺序:将sleep(1)移到循环末尾,确保每次循环先处理消息队列和子进程状态,再进入睡眠,避免消息处理被延迟:
while (1)
{
    // 先处理消息队列请求
    num_bytes = msgrcv(queueID, &mt_msg, MSGUSER_SIZE, MSG_BUDGET, IPC_NOWAIT);
    if (num_bytes >= 0)
    {
        receivedBudgetRequest(mt_msg,queueID);
    }

    // 再处理子进程退出
    pidnow= waitpid(-1, NULL, WNOHANG);
    if(pidnow ==-1)perror("error in waitpid: ");
    if(pidnow > 0){
        // 原有子进程退出处理逻辑
    }

    // 检查是否所有进程已终止
    if (user_done+node_done == SO_USERS_NUM+SO_NODES_NUM) {
        break;
    }

    // 最后执行睡眠
    sleep(1);
}
  1. 强制刷新stdout缓冲区:在user的日志打印后添加fflush(stdout),或者在日志末尾加上换行符\n,确保内容立即输出:
// 方式1:添加换行
printf("\n sono l'user: %d\n",getpid());
// 方式2:手动刷新
printf("\n sono l'user: %d",getpid());
fflush(stdout);
  1. 完善消息发送错误处理:在master的receivedBudgetRequest中添加msgsnd的错误检查,确保回复消息能正确发送:
void receivedBudgetRequest(struct msgUser mt_msg,int queueID)
{
  int initialBudget = calculateBudgetFor(mt_msg.pid);
  struct msgUser messageMaster;
  messageMaster.mtype = mt_msg.pid; 
  messageMaster.pid = getpid();
  messageMaster.numero = initialBudget; 

  if(msgsnd(queueID, &messageMaster, MSGUSER_SIZE, 0) == -1){
      perror("\nerror in msgsnd to user");
  }
  printf("\n mando messaggio ad user\n");
  fflush(stdout);
}

额外说明

  • IPC_NOWAIT的使用是合理的,避免master卡在消息队列上无法处理子进程退出;
  • 若user仍出现异常,可检查calculateBudgetFor函数是否正确返回预算值,以及消息队列的权限(当前0600确保只有当前用户可访问,符合需求)。

内容的提问来源于stack exchange,提问作者Matteo Pagliarello

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 15:25:17