基于线程实现消息队列:输出重叠问题的解决方案探讨
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
相关产品推荐
相关产品推荐

