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

基于线程实现消息队列:输出重叠问题的解决方案探讨

POSIX消息队列线程程序:输出重叠问题解决与代码优化

你实现的基于POSIX消息队列的线程聊天程序,出现server和client线程输出内容重叠的核心原因是**printf是非线程安全函数**——当两个线程同时调用printf时,操作系统的线程调度可能打断输出流程,导致内容交叉。用sleep(0.9)只是靠延迟线程执行碰巧规避问题,属于不可靠的临时方案,一旦系统负载变化,重叠问题大概率会重现。

正确解决方法:互斥锁同步输出

通过引入pthread互斥锁,强制同一时间只有一个线程能执行输出操作,从根源上避免输出冲突。

修改后的完整代码

#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include <pthread.h>
#include <mqueue.h>

#define MAX_SIZE    1024
#define QUEUE_NAME  "/MQ"
#define MSG_STOP    "END"

// 全局互斥锁,用于同步所有标准输出操作
pthread_mutex_t print_mutex;

void *server() {
    mqd_t mq;
    struct mq_attr attr;
    char buffer[MAX_SIZE + 1];
    int state = 1;

    attr.mq_flags = 0;
    attr.mq_maxmsg = 100;
    attr.mq_msgsize = MAX_SIZE;
    attr.mq_curmsgs = 0;

    // 检查消息队列打开结果
    mq = mq_open(QUEUE_NAME, O_CREAT | O_RDONLY, 0644, &attr);
    if (mq == (mqd_t)-1) {
        perror("server: mq_open failed");
        pthread_mutex_lock(&print_mutex);
        fprintf(stderr, "Server线程无法打开消息队列\n");
        pthread_mutex_unlock(&print_mutex);
        return NULL;
    }

    while (state) {
        // 接收消息并检查结果
        ssize_t recv_len = mq_receive(mq, buffer, MAX_SIZE, NULL);
        if (recv_len == -1) {
            perror("server: mq_receive failed");
            break;
        }
        buffer[recv_len] = '\0'; // 确保字符串终止

        // 加锁后执行输出
        pthread_mutex_lock(&print_mutex);
        if (!strncmp(buffer, MSG_STOP, strlen(MSG_STOP))) {
            state = 0;
            printf("chat closing\n");
        } else {
            printf("Msg reached: %s\n", buffer);
        }
        pthread_mutex_unlock(&print_mutex);
    }

    // 关闭消息队列
    if (mq_close(mq) == -1) {
        perror("server: mq_close failed");
    }
    return NULL;
}

void *client() {
    mqd_t mq;
    char buffer[MAX_SIZE];
    int state = 1;

    // 检查消息队列打开结果
    mq = mq_open(QUEUE_NAME, O_WRONLY);
    if (mq == (mqd_t)-1) {
        perror("client: mq_open failed");
        pthread_mutex_lock(&print_mutex);
        fprintf(stderr, "Client线程无法打开消息队列\n");
        pthread_mutex_unlock(&print_mutex);
        return NULL;
    }

    // 加锁输出提示信息
    pthread_mutex_lock(&print_mutex);
    printf("To exit the chat type END \n");
    pthread_mutex_unlock(&print_mutex);

    while (state) {
        // 加锁输出输入提示
        pthread_mutex_lock(&print_mutex);
        printf("chatting: ");
        fflush(stdout);
        pthread_mutex_unlock(&print_mutex);

        memset(buffer, 0, MAX_SIZE);
        // 读取输入并检查结果
        if (fgets(buffer, MAX_SIZE, stdin) == NULL) {
            pthread_mutex_lock(&print_mutex);
            fprintf(stderr, "读取输入失败\n");
            pthread_mutex_unlock(&print_mutex);
            break;
        }

        // 移除fgets读取的换行符
        buffer[strcspn(buffer, "\n")] = '\0';

        if (!strncmp(buffer, MSG_STOP, strlen(MSG_STOP))) {
            // 发送终止消息,只传有效长度
            mq_send(mq, buffer, strlen(buffer)+1, 0);
            state = 0;
        } else {
            // 发送消息并检查结果
            if (mq_send(mq, buffer, strlen(buffer)+1, 0) == -1) {
                perror("client: mq_send failed");
                pthread_mutex_lock(&print_mutex);
                fprintf(stderr, "发送消息失败\n");
                pthread_mutex_unlock(&print_mutex);
                break;
            }

            // 加锁输出发送确认
            pthread_mutex_lock(&print_mutex);
            printf("Msg sent: %s\n", buffer);
            pthread_mutex_unlock(&print_mutex);
        }
    }

    // 关闭消息队列
    if (mq_close(mq) == -1) {
        perror("client: mq_close failed");
    }
    return NULL;
}

int main(int argc, char *argv[]) {
    pthread_t t1, t2;
    int ret;

    // 初始化互斥锁
    if (pthread_mutex_init(&print_mutex, NULL) != 0) {
        perror("pthread_mutex_init failed");
        return EXIT_FAILURE;
    }

    // 创建server线程并检查结果
    ret = pthread_create(&t1, NULL, &server, NULL);
    if (ret != 0) {
        perror("pthread_create server failed");
        pthread_mutex_destroy(&print_mutex);
        return EXIT_FAILURE;
    }

    // 创建client线程并检查结果
    ret = pthread_create(&t2, NULL, &client, NULL);
    if (ret != 0) {
        perror("pthread_create client failed");
        pthread_cancel(t1);
        pthread_join(t1, NULL);
        pthread_mutex_destroy(&print_mutex);
        return EXIT_FAILURE;
    }

    pthread_join(t1, NULL);
    pthread_join(t2, NULL);

    // 销毁互斥锁
    pthread_mutex_destroy(&print_mutex);

    // 程序退出时移除消息队列,避免系统资源泄漏
    if (mq_unlink(QUEUE_NAME) == -1) {
        perror("mq_unlink failed");
    }

    return EXIT_SUCCESS;
}

关键代码优化建议

  • 强制线程安全输出:所有printf、fprintf等标准输出操作都必须被互斥锁包裹,彻底解决输出重叠问题。
  • 完善错误处理:原代码未检查mq_open、mq_receive、pthread_create等函数的返回值,实际开发中必须添加错误处理,避免程序静默崩溃。
  • 优化消息传输长度:原代码固定发送MAX_SIZE字节,修改后只发送实际输入的字符串长度(加1用于终止符),减少不必要的内存开销。
  • 清理系统资源:程序退出时调用mq_close关闭队列、mq_unlink移除队列文件,销毁互斥锁,避免资源泄漏。
  • 移除sleep依赖:完全用互斥锁替代sleep,既保证正确性,又不会浪费CPU时间。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 08:54:56