AsyncQuery实现阻塞问题排查及环形队列方案咨询
AsyncQuery问题分析与解决方案
一、当前实现的阻塞原因分析
- 线程退出逻辑缺陷:析构函数仅设置
_run = false,但工作线程在_cv.wait时的条件仅检查!_messages.empty()。当队列为空时,工作线程会一直阻塞在wait调用上,无法检测到_run的变化,导致_threadCall.join()永久阻塞,测试无法完成。 - 线程安全问题:
Count()函数未对_count的读取加锁,而_count在工作线程中被修改,存在数据竞争,可能导致测试断言结果不符合预期。 - 测试等待时间不足:
AddDelayTest中每个回调需要Sleep(10),10个消息至少需要100ms处理时间,但测试仅Sleep(10),此时工作线程还未处理完所有消息,_count未达到10,断言失败。
二、修正方案
1. 修复线程退出逻辑
修改CallRunner中的wait条件,加入!_run判断,确保线程能在退出信号触发时唤醒;同时析构函数中添加条件变量通知,唤醒阻塞的工作线程:
void AsyncQuery::CallRunner() { while (_run) { std::unique_lock readPopLock(_mtx); // 当队列空且未收到退出信号时等待 _cv.wait(readPopLock, [this] { return !_messages.empty() || !_run; }); if (!_run) break; // 收到退出信号,直接退出循环 auto msg = _messages.front(); _messages.pop(); readPopLock.unlock(); _log(msg.c_str()); std::lock_guard<std::mutex> countLock(_mtx); _count++; } } AsyncQuery::~AsyncQuery() { { std::lock_guard<std::mutex> lock(_mtx); _run = false; } _cv.notify_all(); // 唤醒等待的工作线程 _threadCall.join(); }
2. 修复线程安全的计数读取
给Count()函数添加锁,确保读取_count时的线程安全:
int AsyncQuery::Count() { std::lock_guard<std::mutex> lock(_mtx); return _count; }
3. 调整测试等待时间
修改AddDelayTest中的等待时间,确保所有消息处理完成:
TEST_METHOD(AddDelayTest) { AsyncQuery sut(logDelay); for (int i = 0; i < 10; ++i) { std::wstringstream wstr{}; wstr << L"AddDelay " << i << L'\n'; sut.Add(wstr.str()); } Sleep(110); // 10个消息×10ms + 预留缓冲时间 Assert::AreEqual(10, sut.Count()); }
三、双交替队列方案的合理性分析
1. 方案合理性
采用两个std::queue<std::wstring>交替的环形队列方案是合理的,核心优势在于减少锁的竞争时间:
- 生产者向其中一个空闲队列写入消息,无需等待消费者处理;
- 消费者从另一个满队列读取并处理,同样不干扰生产者;
- 仅当队列切换时(生产者写满当前队列、消费者处理完当前队列),需要短暂加锁完成队列交换。
这种设计非常适配当前场景:生产者快速写入、消费者处理缓慢,能有效避免生产者因锁阻塞,提升消息接收的吞吐量。
2. 是否需要条件变量
仍然需要条件变量。当两个队列都为空时,消费者没有消息可处理,需要进入等待状态;当生产者向空队列写入消息后,需要通过条件变量唤醒消费者。此外,在收到退出信号时,也需要通过条件变量唤醒阻塞的消费者线程,确保线程正常退出。
内容的提问来源于stack exchange,提问作者Soleil
相关产品推荐
相关产品推荐

