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

AMQP-CPP库自定义select事件循环无法运行问题排查

解决AMQP-CPP自定义select事件循环无法触发onConnected的问题

我仔细看了你的代码和问题描述,发现你在自定义select事件循环时犯了几个关键错误,导致连接无法建立,onConnected回调始终不触发。咱们一步步来修复这些问题:

1. 错误的文件描述符监听管理

你的MyTcpHandler只保存了单个m_fd和m_flags,而且在monitor方法里只处理了AMQP::readable的情况,完全忽略了写事件(AMQP::writable),也没处理移除监听的场景(当flags == 0时)。AMQP连接建立过程中需要向服务器发送握手数据,必须监听写事件才能完成初始化流程。

另外,你直接在monitor里修改全局的rfds,但每次循环都会调用FD_ZERO(&rfds)重置这个集合,这会导致之前的监听设置被清空。正确的做法是把需要监听的fd和对应的事件保存在一个容器里,每次循环重新构建fd集合。

2. Select调用未处理写事件

你的select调用只监听了读集合,没有传入写集合,这会导致AMQP-CPP无法发送初始化数据,连接握手永远无法完成。

3. Process方法调用参数错误

你调用connection.process(maxfd, myHandler.m_flags)的参数完全不对,process方法需要的是当前触发的事件(AMQP::readable或AMQP::writable)以及对应的fd,而不是maxfd和保存的flags。

4. Timeval未在循环中重置

select函数会修改传入的timeval结构体,把剩余时间写进去,所以每次循环都需要重新初始化这个结构体,否则后续的select会使用上次剩余的时间,可能导致立即超时。


修正后的完整代码

#include <iostream>
#include <sys/select.h>
#include <cop_amqpcpp.h>
#include <cop_amqpcpp/linux_tcp.h>
#include <unordered_map>
#include <errno.h>
#include <string.h>
#include <algorithm>

using namespace std;

class MyTcpHandler : public AMQP::TcpHandler {
private:
    // 保存需要监听的fd和对应的事件标志
    unordered_map<int, int> _watches;

    virtual void onConnected(AMQP::TcpConnection* connection) {
        std::cout << "connected" << std::endl;
    }

    virtual void onError(AMQP::TcpConnection* connection, const char* message) {
        std::cout << "AMQP TCPConnection error: " << message << std::endl;
    }

    virtual void onClosed(AMQP::TcpConnection* connection) {
        std::cout << "closed" << std::endl;
    }

    virtual void monitor(AMQP::TcpConnection* connection, int fd, int flags) {
        if (flags == 0) {
            // 移除监听的fd
            _watches.erase(fd);
            return;
        }
        cout << "Fd " << fd << " Flags " << flags << endl;
        // 更新或添加监听事件
        _watches[fd] = flags;
    }

public:
    // 对外暴露监听列表
    const unordered_map<int, int>& watches() const {
        return _watches;
    }
};

int main(int argc, char* argv[]) {
    MyTcpHandler myHandler;
    int maxfd = 0;
    int result = 0;

    // 服务器地址
    AMQP::Address address("amqp://test:test@localhost/");

    // 创建连接和通道
    AMQP::TcpConnection connection(&myHandler, address);
    AMQP::TcpChannel channel(&connection);
    
    channel.onError([](const char* message) {
        cout << "channel error: " << message << endl;
    });
    channel.onReady([]() {
        cout << "channel ready " << endl;
    });

    // 声明交换机、队列并绑定
    channel.declareExchange("my-exchange", AMQP::fanout);
    channel.declareQueue("my-queue");
    channel.bindQueue("my-exchange", "my-queue", "my-routing-key");

    do {
        fd_set rfds, wfds;
        FD_ZERO(&rfds);
        FD_ZERO(&wfds);
        maxfd = 0;

        // 添加标准输入到监听(可选,用于手动退出程序)
        FD_SET(0, &rfds);
        maxfd = max(maxfd, 0);

        // 遍历所有需要监听的fd,分别添加到读/写集合
        const auto& watches = myHandler.watches();
        for (const auto& pair : watches) {
            int fd = pair.first;
            int flags = pair.second;

            if (flags & AMQP::readable) FD_SET(fd, &rfds);
            if (flags & AMQP::writable) FD_SET(fd, &wfds);
            maxfd = max(maxfd, fd);
        }

        // 每次循环重置超时时间,避免select修改后影响下一次调用
        struct timeval tv { 1, 0 };

        result = select(maxfd + 1, &rfds, &wfds, NULL, &tv);

        if (result == -1) {
            if (errno == EINTR) {
                // 被信号中断,继续循环即可
                continue;
            }
            cout << "Select error: " << strerror(errno) << endl;
            break;
        } else if (result > 0) {
            // 处理每个触发事件的fd
            const auto& watches = myHandler.watches();
            for (const auto& pair : watches) {
                int fd = pair.first;
                int triggerFlags = 0;

                if (FD_ISSET(fd, &rfds)) triggerFlags |= AMQP::readable;
                if (FD_ISSET(fd, &wfds)) triggerFlags |= AMQP::writable;

                if (triggerFlags != 0) {
                    connection.process(fd, triggerFlags);
                }
            }

            // 处理标准输入(输入任意内容退出程序)
            if (FD_ISSET(0, &rfds)) {
                char buf[1024];
                read(0, buf, sizeof(buf));
                cout << "Exiting..." << endl;
                break;
            }
        }

    } while (1);

    return 0;
}

关键修正说明

  • 用unordered_map<int, int>存储所有需要监听的fd和事件,支持动态添加/移除监听
  • monitor方法现在正确处理添加、更新和移除监听的逻辑
  • 每次循环重新构建读/写fd集合,确保所有需要监听的fd都被正确加入
  • Select调用同时传入读集合和写集合,支持AMQP的写操作需求
  • connection.process传入正确的触发fd和事件标志,符合AMQP-CPP的调用规范
  • 每次循环重置timeval结构体,避免超时时间被修改后导致的异常
  • 增加了更完善的错误处理逻辑

内容的提问来源于stack exchange,提问作者selbolder

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:16:55