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
相关产品推荐
相关产品推荐

