线程退出时无锁MPMC队列出现消息丢失问题求助
问题描述
我在实现互斥算法时使用无锁MPMC队列,遇到了消息丢失的不一致行为,具体现象如下:
- 终止流程开始前无消息丢失
- 若临界区代码极简(仅短毫秒级sleep),终止流程会导致消息丢失,线程因等待套接字挂起
- 若临界区包含实际工作(如向无关套接字发送TCP请求),终止流程可成功完成,所有节点优雅停机
相关代码与架构
我难以复现最小示例,但基础架构如下:
- N个节点分布在N台不同服务器上
- 每个节点包含:
- 应用线程(执行临界区)
- N-1个读取线程(接收其他节点消息)
- 协调线程(实现互斥协议)
- 写入线程(向其他节点发送消息)
消息传递路径:
- 读取线程 -> MPMC队列 -> 协调线程(来自其他节点的消息)
- 协调线程 -> MPMC队列 -> 写入线程(待发送给其他节点的消息)
- 应用线程 -> MPMC队列 -> 协调线程(请求/退出临界区)
以下是我认为相关的代码片段:
应用线程
for (int i = 0; i < 10; i++) { client.cs_enter() // Blocks until ready // do_some_work() sleep(5) client.cs_leave() } // wait for threads
消息结构
enum MessageKind { REQUEST, REPLY, CS_ENTER, CS_LEAVE, DONE, TERMINATE }; struct Message { MessageKind kind; uint32_t node; uint32_t timestamp; };
写入线程
while (true) { std::pair<uint32_t, Message> tm({ 0, Message::cs_enter() }); q->wait_dequeue(tm); // MAX node means no more messages to send if (tm.first == std::numeric_limits<uint32_t>::max()) break; // map contains node_id -> socket_fd mapping send_message((*map)[tm.first], tm.second); } // exit thread
读取线程
while (true) { Message m = recv_message(fd); MessageKind k = m.kind; q->enqueue(m); if (k == TERMINATE) break; } // exit thread
协调线程
while(true) { if (this->done && this->terminated.size() == this->config->nodes.size() - 1) break; Message m = Message::cs_enter(); this->input_queue->wait_dequeue(m); switch (m.kind) { // handle message based on kind } } // inform writer no more work to do this->output_queue->enqueue({ std::numeric_limits<uint32_t>::max(), Message::done() }); // exit thread
我完全不知道该如何推进。随意增加临界区执行时间是不可靠的解决方案,也不确定它在所有测试环境中的表现。我也不理解为何增加临界区执行时间能解决问题,因为它与协议代码无关。
我的代码语义是否无法保证消息交付?更具体地说,以下代码能否保证消费者收到消息?
// Producer - thread 1 q->enqueue(message); thread_exit() // Consumer - thread 2 Message m; q->wait_dequeue(m);
还是需要额外编写代码(如内存屏障等)来保证该特性?
内容的提问来源于stack exchange,提问作者Kungfunk
相关产品推荐
相关产品推荐

