使用libpqxx订阅两个通知通道无法接收消息求助
问题分析与解决
你的代码卡在循环里收不到通知,核心问题是libpqxx不会自动处理连接事件——空循环while(!L.done()){}只是占用CPU空转,没有触发连接去读取并处理PostgreSQL发送的通知消息,导致notification_receiver的回调函数永远不会被调用。
问题代码回顾
TestListener类
class TestListener final : public pqxx::notification_receiver { bool m_done; public: explicit TestListener(pqxx::connection &conn, std::string Name) : pqxx::notification_receiver(conn, Name), m_done(false) {} void operator()(std::string const &, int be_pid) override { m_done = true; std::cout << "Received notification: " << channel() << " pid=" << be_pid << std::endl; } bool done() const { return m_done; } };
等待函数与主函数
void waitForChannel(TestListener &L){ while(!L.done()){}; std::cout << "received notificaton for channel " << L.channel() << "\n"; }; int main(){ // assume the connection is made // ... TestListener L1{conn, "channel1"}; TestListener L2{conn, "channel2"}; pqxx::perform([&conn, &L1] { pqxx::work tx{conn}; tx.exec0("NOTIFY " + tx.quote_name(L1.channel())); tx.commit(); }); pqxx::perform([&conn, &L2] { pqxx::work tx{conn}; tx.exec0("NOTIFY " + tx.quote_name(L2.channel())); tx.commit(); }); while(true){ waitForChannel(L1); waitForChannel(L2); } return 0; }
修复方案
1. 修改等待函数,让连接主动处理事件
libpqxx的connection::wait()方法会阻塞等待连接上的事件(包括通知),并触发对应的notification_receiver回调。修改waitForChannel函数:
void waitForChannel(pqxx::connection &conn, TestListener &L){ while(!L.done()){ conn.wait(); // 关键:让连接处理事件,触发通知回调 }; std::cout << "收到通道 " << L.channel() << " 的通知\n"; };
2. 调整主函数的调用逻辑
调用等待函数时传入连接对象,同时添加重置状态的方法(否则第一次收到通知后,后续循环会直接跳过):
先给TestListener类添加重置方法:
void reset() { m_done = false; }
再修改主函数的循环:
while(true){ waitForChannel(conn, L1); L1.reset(); // 重置状态,等待下一次通知 waitForChannel(conn, L2); L2.reset(); // 重置状态,等待下一次通知 }
3. 可选:非阻塞方式同时监听多通道
如果需要同时监听多个通道,避免阻塞在单个通道的等待上,可以使用connection::poll()配合短睡眠:
#include <algorithm> #include <chrono> #include <thread> void waitForAnyNotification(pqxx::connection &conn, std::vector<TestListener*> listeners){ bool any_done = false; while(!any_done){ conn.poll(); // 非阻塞处理事件 std::this_thread::sleep_for(std::chrono::milliseconds(100)); // 避免CPU空转 any_done = std::any_of(listeners.begin(), listeners.end(), [](TestListener* l){ return l->done(); }); } }
核心原理
PostgreSQL的NOTIFY消息通过数据库连接发送,libpqxx需要主动从连接读取这些消息,才能触发notification_receiver的回调函数。conn.wait()或conn.poll()是让连接处理底层消息的关键操作,原代码缺少这一步,导致回调永远不会执行。
内容的提问来源于stack exchange,提问作者shashashamti2008
相关产品推荐
相关产品推荐

