libpqxx await_notification()误触发信号问题排查求助
问题
使用libpqxx C++库监听TimescaleDB的两个触发器信号,通过两个线程分别执行无限循环等待触发:
while (true){ conn.await_notification(); spdlog::info("channel1_updated"); }
while (true){ conn.await_notification(); spdlog::info("channel2_updated"); }
但未生成任何触发器信号时,channel1_updated和channel2_updated日志被自动打印,conn.await_notification()未按预期阻塞代码。完整代码如下:
#include "pqxx/config-public-compiler.h" #include <cctype> #include <cerrno> #include <cstring> #include <ctime> #include <iostream> #include <pqxx/internal/header-pre.hxx> #include <pqxx/internal/wait.hxx> #include <pqxx/internal/header-post.hxx> #include <pqxx/notification> #include <pqxx/transaction> #include <pqxx/transactor> #include <thread> #include "spdlog/spdlog.h" #include "spdlog/cfg/env.h" #include "spdlog/fmt/ostr.h" class PostgresChannelListener final : public pqxx::notification_receiver { bool m_done; public: int backend_pid; explicit PostgresChannelListener(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; this->backend_pid = be_pid; spdlog::info("Received notification: {} pid=", channel(), be_pid); } bool done() const { return m_done; } }; int waitForChannelUpdated(const std::string& channel){ try { pqxx::connection conn("dbname=postgres user=postgres password=password hostaddr=127.0.0.1 port=5432"); if (conn.is_open()) { spdlog::info("Opened database successfully: {}", conn.dbname()); spdlog::info("PostgreSQL server version: {}", conn.server_version()); } else { spdlog::error("Failed to open database"); return 1; } PostgresChannelListener L{conn, channel}; pqxx::perform([&conn, &L] { pqxx::work tx{conn}; tx.exec0("NOTIFY " + tx.quote_name(L.channel())); tx.commit(); }); while (true) { conn.await_notification(); spdlog::info(channel); } conn.close(); spdlog::info("Closed database successfully"); } catch (const std::exception &e) { spdlog::error("Error: {}", e.what()); return 1; } return 0; } int main() { std::thread threadChannel1([](){ waitForChannelUpdated("channel1_updated"); }); std::thread threadChannel2([](){ waitForChannelUpdated("channel2_updated"); }); threadChannel1.join(); threadChannel2.join(); return 0; }
分析与解决方案
问题根源
代码启动后主动向监听频道发送了NOTIFY指令,导致await_notification()被直接唤醒:
在waitForChannelUpdated函数中,创建PostgresChannelListener后执行的这段代码,会主动给当前订阅的频道发一条通知,线程启动后立刻收到这条自发送的信号,触发日志打印:
pqxx::perform([&conn, &L] { pqxx::work tx{conn}; tx.exec0("NOTIFY " + tx.quote_name(L.channel())); tx.commit(); });
修复步骤
- 移除主动发送NOTIFY的代码块:触发器会在数据变更时自动发送信号,无需手动触发。
- 匹配目标频道通知:直接使用
await_notification()时,需要判断收到的通知是否属于当前线程监听的频道,避免处理其他频道的信号。
修复后的示例代码:
int waitForChannelUpdated(const std::string& channel){ try { pqxx::connection conn("dbname=postgres user=postgres password=password hostaddr=127.0.0.1 port=5432"); if (conn.is_open()) { spdlog::info("Opened database successfully: {}", conn.dbname()); spdlog::info("PostgreSQL server version: {}", conn.server_version()); } else { spdlog::error("Failed to open database"); return 1; } // 订阅目标频道(使用notification_receiver时会自动执行LISTEN,此处如果保留Listener可省略此段) pqxx::work tx{conn}; tx.exec0("LISTEN " + tx.quote_name(channel)); tx.commit(); while (true) { // 获取通知详情 auto [notify_channel, be_pid, payload] = conn.await_notification(); // 仅处理当前线程监听频道的通知 if (notify_channel == channel) { spdlog::info(channel); } } conn.close(); spdlog::info("Closed database successfully"); } catch (const std::exception &e) { spdlog::error("Error: {}", e.what()); return 1; } return 0; }
补充说明
如果要使用PostgresChannelListener的回调机制,需调整循环逻辑:
while (true) { // 处理到达的通知,触发回调函数 conn.process_notifications(); // 添加短暂休眠避免CPU占用过高 std::this_thread::sleep_for(std::chrono::milliseconds(100)); }
内容的提问来源于stack exchange,提问作者shashashamti2008
相关产品推荐
相关产品推荐

