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

TCP Socket连接成功但输出流无法写入,未收到预期回调求助

TCP Socket连接成功但写入无响应的问题排查与修复

TCP Socket连接已成功建立,但无法向输出流写入任何数据。通常客户端向服务器成功发送消息后,代理会触发bytesAvailable回调;当前已收到SpaceAvailable回调且日志显示消息已发送,但未收到预期的回调。客户端本身无问题,旧Objective-C代码运行正常。相关Swift代码如下:

class ViewController: UIViewController, StreamDelegate {
    var inputStream: InputStream?
    var outputStream: OutputStream?
    
    var messageQueue = [String]()
    
    var timer : Timer!
    var writeReady = false
    
    var timeInterval = 2.0
    
    var clientSocket: Int32!
    
    override func viewDidLoad() {
        super.viewDidLoad()
        
        self.startTCP()
    }
    
    func startTCP(){
        self.connectToServer(host: "192.168.0.200", port: 50001)
    }
    
    func connectToServer(host: String, port: Int) {
        Stream.getStreamsToHost(withName: host, port: port, inputStream: &inputStream, outputStream: &outputStream)
        
        inputStream?.delegate = self
        outputStream?.delegate = self
        
        inputStream?.schedule(in: RunLoop.current, forMode: .common)
        outputStream?.schedule(in: RunLoop.current, forMode: .common)
        
        inputStream?.open()
        outputStream?.open()
        print("connection successful")
    }
    
    func stream(_ aStream: Stream, handle eventCode: Stream.Event) {
        switch eventCode {
        case .openCompleted:
            if aStream == outputStream {
                print("open completed for output stream")
                writeReady = true
                sendConnectionMessage()
                startHeartBeat()
            }
        case .hasBytesAvailable:
            print("bytes avaialble")
            guard let inputStream = inputStream else {
                print("Input stream is nil")
                return
            }
            let bufferSize = 1024
            var buffer = Array<UInt8>(repeating: 0, count: bufferSize)
            let bytesRead = inputStream.read(&buffer, maxLength: bufferSize)
            if bytesRead > 0 {
                let receivedMessage = String(bytes: buffer, encoding: .utf8)!
                print("Received message: \(receivedMessage)")
            }
        case .errorOccurred:
            print("Error occurred: \(aStream.streamError!.localizedDescription)")
            disconnect()
            stopHeartBeat()
            startTCP()
        case .endEncountered:
            print("End encountered")
        case .hasSpaceAvailable:
            print("space available")
            
            if(aStream == self.outputStream)
            {
                handleMessageQueue()
            }
        default:
            print("Unknown event")
        }
    }
    
    func handleMessageQueue(){
        if messageQueue.count>0 {
            if let message = messageQueue.first {
                sendMessage(message)
                messageQueue.removeFirst()
            }
        } else {
            //handle empty queue
        }
    }
    
    @objc func sendConnectionMessage() {
        messageQueue.append("{\"msg.type\":\"Connection\",\"source.address\":\"192.168.0.126\"}")
    }
    
    func sendMessage(_ message: String) {
        if !writeReady{ return }
        
        guard let outputStream = outputStream else {
            print("Output stream is nil")
            return
        }

        let data = message.data(using: .utf8)!
        print("writing")
        //        let bytesWritten = outputStream.write(data: data)
        
        let bytesWritten = data.withUnsafeBytes { outputStream.write($0, maxLength: data.count) }
        
        if bytesWritten < 0 {
            print("Error sending message")
        } else {
            print("Sent message: \(message)")
        }
    }
    
    func startHeartBeat(){
        timer = Timer.scheduledTimer(withTimeInterval: timeInterval, repeats: true, block: { _ in
            self.sendMessage("{\"msg.data\":[{\"Dev.ID\":\"0\",\"Dev.inst\":\"0\"}],\"msg.type\":\"Register\"}\n")
        })
    }
    
    func stopHeartBeat(){
        if timer == nil { return }
        timer.invalidate()
    }
    
    func disconnect() {
        inputStream?.close()
        outputStream?.close()
        inputStream = nil
        outputStream = nil
    }
}

核心问题分析

  • 部分写入未处理:TCP是面向流的协议,OutputStream.write()可能无法一次性写入全部数据,当前代码仅调用一次write,剩余数据会丢失,导致服务器无法收到完整消息。
  • 心跳包未走消息队列:心跳定时器直接调用sendMessage,跳过了队列机制,可能在输出流状态不稳定时发送失败,或者打乱消息发送顺序。
  • 读取数据处理不严谨:直接将缓冲区转为String,未考虑数据截断(比如消息超过1024字节)或编码错误的情况,导致无法正确解析服务器响应。
  • 可写状态管理缺失:仅在openCompleted时设置writeReady = true,但hasSpaceAvailable触发后应该确保可写状态持续有效,避免后续写入被拦截。

修复后的代码

class ViewController: UIViewController, StreamDelegate {
    var inputStream: InputStream?
    var outputStream: OutputStream?
    
    var messageQueue = [Data]() // 改用Data存储,避免重复转码
    var currentWritingData: Data?
    var currentWriteOffset = 0
    
    var timer: Timer? // 改用可选类型,避免强制解包
    var writeReady = false
    
    let timeInterval = 2.0
    
    override func viewDidLoad() {
        super.viewDidLoad()
        startTCP()
    }
    
    func startTCP(){
        connectToServer(host: "192.168.0.200", port: 50001)
    }
    
    func connectToServer(host: String, port: Int) {
        Stream.getStreamsToHost(withName: host, port: port, inputStream: &inputStream, outputStream: &outputStream)
        
        inputStream?.delegate = self
        outputStream?.delegate = self
        
        inputStream?.schedule(in: RunLoop.current, forMode: .common)
        outputStream?.schedule(in: RunLoop.current, forMode: .common)
        
        inputStream?.open()
        outputStream?.open()
        print("connection successful")
    }
    
    func stream(_ aStream: Stream, handle eventCode: Stream.Event) {
        switch eventCode {
        case .openCompleted:
            if aStream == outputStream {
                print("open completed for output stream")
                writeReady = true
                sendConnectionMessage()
                startHeartBeat()
            }
        case .hasBytesAvailable:
            print("bytes available")
            guard let inputStream = inputStream else {
                print("Input stream is nil")
                return
            }
            let bufferSize = 1024
            var buffer = Array<UInt8>(repeating: 0, count: bufferSize)
            var receivedData = Data()
            
            while inputStream.hasBytesAvailable {
                let bytesRead = inputStream.read(&buffer, maxLength: bufferSize)
                if bytesRead > 0 {
                    receivedData.append(buffer, count: bytesRead)
                } else if bytesRead < 0 {
                    print("Read error: \(inputStream.streamError?.localizedDescription ?? "unknown")")
                    break
                }
            }
            
            if !receivedData.isEmpty, let receivedMessage = String(data: receivedData, encoding: .utf8) {
                print("Received message: \(receivedMessage)")
            }
        case .errorOccurred:
            print("Error occurred: \(aStream.streamError!.localizedDescription)")
            disconnect()
            stopHeartBeat()
            startTCP()
        case .endEncountered:
            print("End encountered")
            disconnect()
            stopHeartBeat()
        case .hasSpaceAvailable:
            print("space available")
            writeReady = true // 触发此事件时标记可写
            if aStream == outputStream {
                handleMessageQueue()
            }
        default:
            print("Unknown event")
        }
    }
    
    func handleMessageQueue(){
        guard writeReady, let outputStream = outputStream else { return }
        
        // 先处理当前未写完的数据
        if let data = currentWritingData {
            let remainingBytes = data.count - currentWriteOffset
            let bytesWritten = data.withUnsafeBytes {
                outputStream.write($0.baseAddress!.advanced(by: currentWriteOffset), maxLength: remainingBytes)
            }
            
            if bytesWritten > 0 {
                currentWriteOffset += bytesWritten
                if currentWriteOffset == data.count {
                    // 当前数据写完,清空状态
                    currentWritingData = nil
                    currentWriteOffset = 0
                } else {
                    // 还有剩余数据,等待下一次hasSpaceAvailable
                    return
                }
            } else if bytesWritten < 0 {
                print("Error writing remaining data")
                currentWritingData = nil
                currentWriteOffset = 0
            }
        }
        
        // 处理队列中的下一条消息
        if !messageQueue.isEmpty, currentWritingData == nil {
            currentWritingData = messageQueue.removeFirst()
            handleMessageQueue() // 递归调用,开始写入新数据
        }
    }
    
    func sendConnectionMessage() {
        let message = "{\"msg.type\":\"Connection\",\"source.address\":\"192.168.0.126\"}"
        if let data = message.data(using: .utf8) {
            messageQueue.append(data)
            handleMessageQueue() // 尝试立即写入
        }
    }
    
    func sendMessage(_ message: String) {
        if let data = message.data(using: .utf8) {
            messageQueue.append(data)
            handleMessageQueue() // 尝试立即写入
        }
    }
    
    func startHeartBeat(){
        stopHeartBeat() // 先停止已有定时器,避免重复创建
        timer = Timer.scheduledTimer(withTimeInterval: timeInterval, repeats: true) { [weak self] _ in
            self?.sendMessage("{\"msg.data\":[{\"Dev.ID\":\"0\",\"Dev.inst\":\"0\"}],\"msg.type\":\"Register\"}\n")
        }
    }
    
    func stopHeartBeat(){
        timer?.invalidate()
        timer = nil
    }
    
    func disconnect() {
        inputStream?.close()
        outputStream?.close()
        inputStream = nil
        outputStream = nil
        messageQueue.removeAll()
        currentWritingData = nil
        currentWriteOffset = 0
        writeReady = false
    }
}

关键修改点说明

  • 消息队列改为存储Data:避免重复进行String到Data的转码,提升效率。
  • 处理部分写入:增加currentWritingData和currentWriteOffset跟踪未写完的数据,确保所有字节都能发送到服务器。
  • 心跳包加入队列:心跳消息通过sendMessage加入队列,统一由handleMessageQueue管理发送顺序和状态。
  • 优化数据读取:循环读取所有可用字节,拼接成完整的Data后再转String,避免截断问题。
  • 修正可写状态:在hasSpaceAvailable事件中设置writeReady = true,确保每次可写时都能触发写入。
  • 定时器改为可选类型:避免强制解包导致的崩溃,创建前先停止已有定时器。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 06:32:13