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

iOS+Swift集成RabbitMQ遇错求助:代码排查与解决方案

iOS Swift 集成RabbitMQ通信报错排查与修复

问题场景

在iOS设备上使用Swift通过RabbitMQ与服务器通信时出现连接异常,以下是实现代码及报错日志:

实现代码

import RMQClient

class ViewController: UIViewController {
    
    override func viewDidLoad() {
        super.viewDidLoad()
        // 调用从服务器接收日志的函数
        self.emitLog()
        self.receiveLogs()
        
        // 在指定时机发送消息(如按钮点击等)
       
    }
    
    func createConnection() -> RMQConnection {
        let hostName = "112.219.138.test" // 主机名
        let userName = "test" // 用户名
        let password = "test" // 密码
        let virtualHost = "/" // 虚拟主机(使用默认值"/"请不要修改)

        let uri = "amqp://\(userName):\(password)@\(hostName):5672\(virtualHost)"
        let delegate = RMQConnectionDelegateLogger()
        let connection = RMQConnection(uri: uri, delegate: delegate)
        return connection
    }
    
    func emitLog() {
        let conn = createConnection()
        conn.start()
        let ch = conn.createChannel()
        let x = ch.fanout("logs")
        let msg = "Hello, RabbitMQ!" // 设置要发送的消息内容
        x.publish(msg.data(using: String.Encoding.utf8)!) // 发布消息
        print("Sent \(msg)")
        conn.close()
    }

    
    func receiveLogs() {
        let conn = createConnection()
        conn.start()
        let ch = conn.createChannel()
        let x = ch.fanout("logs")
        let q = ch.queue("", options: .exclusive)
        q.bind(x)
        print("Waiting for logs...")
        q.subscribe({ (_ message: RMQMessage) -> Void in
            if let messageString = String(data: message.body, encoding: .utf8) {
                // 打印接收到的日志
                print("Received Log: \(messageString)")
                
            }
        })
    }
}

报错日志(已翻译)

Sent Hello, RabbitMQ!
Received connection: <RMQConnection: 0x60000301c090> failedToConnectWithError: Error Domain=com.rabbitmq.rabbitmq-objc-client Code=1 "尝试关闭一个已关闭(或从未成功建立)的连接" UserInfo={NSLocalizedDescription=尝试关闭一个已关闭(或从未成功建立)的连接}
Waiting for logs...
线程性能检查器:运行在用户发起级别的QoS线程正在等待运行在后台级别的QoS线程。请研究避免优先级反转的方法
PID: 27587, TID: 8236637
堆栈跟踪

=================================================================

3   RMQClient                           0x0000000101d816cc -[RMQSemaphoreWaiter timesOut] + 92
4   RMQClient                           0x0000000101d3dd8c __23-[RMQConnection start:]_block_invoke + 604
5   libdispatch.dylib                   0x0000000102433747 _dispatch_call_block_and_release + 12
6   libdispatch.dylib                   0x00000001024349f7 _dispatch_client_callout + 8
7   libdispatch.dylib                   0x000000010243c8c9 _dispatch_lane_serial_drain + 1127
8   libdispatch.dylib                   0x000000010243d665 _dispatch_lane_invoke + 441
9   libdispatch.dylib                   0x000000010244a76e _dispatch_root_queue_drain_deferred_wlh + 318
10  libdispatch.dylib                   0x0000000102449b69 _dispatch_workloop_worker_thread + 590
11  libsystem_pthread.dylib             0x0000000101143c47 _pthread_wqthread + 327
12  libsystem_pthread.dylib             0x0000000101142b97 start_wqthread + 15

nw_socket_handle_socket_event [C2:1] Socket SO_ERROR [54: 连接被对等方重置]
nw_socket_handle_socket_event [C1:1] Socket SO_ERROR [54: 连接被对等方重置]
Received connection: <RMQConnection: 0x600003014120> disconnectedWithError: Error Domain=GCDAsyncSocketErrorDomain Code=7 "套接字被远程对等方关闭" UserInfo={NSLocalizedDescription=套接字被远程对等方关闭}
Will start recovery for connection: <RMQConnection: 0x600003014120>
Received connection: <RMQConnection: 0x60000301c090> disconnectedWithError: Error Domain=GCDAsyncSocketErrorDomain Code=7 "套接字被远程对等方关闭" UserInfo={NSLocalizedDescription=套接字被远程对等方关闭}
Starting recovery for connection: <RMQConnection: 0x600003014120>
Recovered connection: <RMQConnection: 0x600003014120>
nw_socket_handle_socket_event [C3:1] Socket SO_ERROR [54: 连接被对等方重置]
Received connection: <RMQConnection: 0x600003014120> disconnectedWithError: Error Domain=GCDAsyncSocketErrorDomain Code=7 "套接字被远程对等方关闭" UserInfo={NSLocalizedDescription=套接字被远程对等方关闭}

错误原因分析

  • 连接生命周期管理错误:emitLog中刚调用publish就立即关闭连接,此时消息可能还未完成发送,直接关闭会导致连接异常,甚至触发服务器端重置连接。
  • 连接未保持强引用:receiveLogs中创建的连接是局部变量,方法执行完毕后会被ARC释放,导致连接被意外断开,触发重连逻辑但最终失败。
  • 线程优先级反转:RMQClient内部使用后台QoS线程,而主线程是用户发起级别,主线程等待后台线程时出现优先级反转,影响连接建立稳定性。
  • 未等待连接就绪:conn.start()是异步操作,直接后续创建通道、发布消息时,连接可能还未成功建立,导致操作无效。

修复方案与代码实现

核心修复点

  1. 复用单个RabbitMQ连接,避免重复创建
  2. 对连接、通道、交换机保持类级别强引用,防止被释放
  3. 等待连接就绪后再执行消息发送/订阅操作
  4. 调整连接的线程QoS,解决优先级反转问题
  5. 合理管理连接关闭时机(仅在不需要通信时关闭)

修复后的完整代码

import RMQClient

class ViewController: UIViewController {
    
    // 保持连接、通道、交换机的强引用
    private var connection: RMQConnection?
    private var channel: RMQChannel?
    private var logsExchange: RMQExchange?
    
    override func viewDidLoad() {
        super.viewDidLoad()
        
        // 初始化RabbitMQ连接
        setupRabbitMQConnection()
    }
    
    private func setupRabbitMQConnection() {
        let hostName = "112.219.138.test"
        let userName = "test"
        let password = "test"
        let virtualHost = "/"
        
        let uri = "amqp://\(userName):\(password)@\(hostName):5672\(virtualHost)"
        let delegate = RMQConnectionDelegateLogger()
        
        // 创建连接时指定线程QoS为用户发起级别,解决优先级反转
        connection = RMQConnection(uri: uri, delegate: delegate, dispatchQueue: DispatchQueue.global(qos: .userInitiated))
        
        // 异步等待连接就绪
        connection?.start { [weak self] in
            guard let self = self, let conn = self.connection else { return }
            
            // 连接就绪后创建通道和交换机
            self.channel = conn.createChannel()
            self.logsExchange = self.channel?.fanout("logs")
            
            // 启动消息订阅
            self.startReceivingLogs()
            
            // 示例:延迟发送消息,确保连接完全就绪
            DispatchQueue.main.asyncAfter(deadline: .now() + 1) {
                self.sendLog(message: "Hello, RabbitMQ!")
            }
        }
    }
    
    private func sendLog(message: String) {
        guard let exchange = logsExchange else {
            print("Exchange未初始化,无法发送消息")
            return
        }
        
        if let data = message.data(using: .utf8) {
            exchange.publish(data)
            print("已发送消息: \(message)")
        }
    }
    
    private func startReceivingLogs() {
        guard let channel = channel, let exchange = logsExchange else {
            print("通道或交换机未初始化,无法订阅消息")
            return
        }
        
        // 创建排他队列并绑定到交换机
        let queue = channel.queue("", options: .exclusive)
        queue.bind(exchange)
        
        print("等待接收日志...")
        
        // 订阅消息
        queue.subscribe { message in
            if let messageString = String(data: message.body, encoding: .utf8) {
                print("收到日志: \(messageString)")
            }
        }
    }
    
    // 在页面销毁时关闭连接
    deinit {
        connection?.close()
        print("RabbitMQ连接已关闭")
    }
}

额外注意事项

  • 确保服务器端RabbitMQ服务正常运行,端口5672对外开放,且用户名/密码、虚拟主机配置正确
  • 网络环境需允许iOS设备访问RabbitMQ服务器(注意防火墙、VPN限制)
  • 如需处理连接断开重连,可通过自定义RMQConnectionDelegate实现重连逻辑

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 09:35:00