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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 12:23:14