Swift 4中InputStream读取TCP消息时barrier使用疑问及崩溃问题
问题分析与解决方案
首先,你确实错误使用了barrier,但这不是导致崩溃的核心原因——真正的问题是你的代码没有维护输入流的读取状态,加上对DispatchQueue barrier的作用理解有误。
1. 为什么Barrier在这里没用?
你创建的inputStreamAccessQueue是串行队列(默认DispatchQueue(label:)创建的是串行队列)。Barrier标记只在并发队列中有用:它用来让某个任务独占队列,确保在它执行前后的其他任务都不会并行运行。但在串行队列中,所有任务本来就是按顺序执行的,加.barrier完全没有效果,不会改变任何行为。
2. 崩溃的真正原因:缺少读取状态管理
InputStream的hasBytesAvailable事件会多次触发——只要流中有可读取的字节就会触发,不管你是否正在处理上一次的读取请求。你的当前代码假设每次handleInput调用都能处理一个完整的消息(先读4字节长度,再读完整内容),但实际情况是:
- 第一次触发时,你读了4字节的长度,然后开始读取消息内容,但还没读完所有
bytes_expected字节时,InputStream又触发了hasBytesAvailable事件; - 这时候
handleInput被再次调用,又去读取4字节的长度,但此时流中剩下的是上一条消息未读完的内容,不是新消息的长度,所以读出来的是垃圾数据,导致bytes_expected变成随机值,最终引发非法内存分配崩溃。
3. 修复方案
你需要维护一个读取状态,记录当前是在读取消息长度,还是在读取消息内容,以及已经读取了多少内容。同时,正确使用串行队列来保护这些状态和流的访问。
修正后的代码示例
private let inputStreamAccessQueue = DispatchQueue(label: "SynchronizedInputStreamAccess") // 维护读取状态,跟踪当前读取阶段 private enum ReadState { case waitingForLength case readingMessage(remainingBytes: Int, currentMessage: Data) } private var readState: ReadState = .waitingForLength func inputStreamHandler(_ event: Stream.Event) { switch event { case Stream.Event.hasBytesAvailable: // 用async提交任务,避免阻塞InputStream所在的RunLoop inputStreamAccessQueue.async { self.handleInput() } // 处理其他流事件 case Stream.Event.endEncountered: inputStreamAccessQueue.async { self.cleanupInputStream() } case Stream.Event.errorOccurred: inputStreamAccessQueue.async { self.errorHandler(NetworkingError.InputError("Stream error occurred")) } default: break } } func handleInput() { guard let istr = self.inputStream, istr.hasBytesAvailable else { log.error(self.buildLogMessage("InputStream unavailable or no bytes available")) return } switch readState { case .waitingForLength: // 读取4字节的消息长度 var lengthBuffer = [UInt8](repeating: 0, count: 4) let bytesRead = istr.read(&lengthBuffer, maxLength: 4) guard bytesRead == 4 else { self.errorHandler(NetworkingError.InputError("Expected 4 bytes for length, got \(bytesRead)")) readState = .waitingForLength // 重置状态 return } // 转换为大端序的UInt32,再转为Int let bytesExpected = Int(UInt32(bigEndian: lengthBuffer.withUnsafeBytes { $0.load(as: UInt32.self) })) log.info(self.buildLogMessage("Expect \(bytesExpected) bytes for message")) // 切换状态到读取消息内容,初始化空Data累积内容 readState = .readingMessage(remainingBytes: bytesExpected, currentMessage: Data()) // 立即尝试读取消息内容(流中可能已有数据) var remaining = bytesExpected var messageData = Data() continueReadingMessage(remainingBytes: &remaining, currentMessage: &messageData) case .readingMessage(var remainingBytes, var currentMessage): continueReadingMessage(remainingBytes: &remainingBytes, currentMessage: ¤tMessage) } } private func continueReadingMessage(remainingBytes: inout Int, currentMessage: inout Data) { guard let istr = self.inputStream, istr.hasBytesAvailable else { // 更新状态,等待下一次流事件触发 readState = .readingMessage(remainingBytes: remainingBytes, currentMessage: currentMessage) return } // 用合理的缓冲区大小(比如1024字节),避免分配过大内存 let bufferSize = min(remainingBytes, 1024) var buffer = [UInt8](repeating: 0, count: bufferSize) let bytesRead = istr.read(&buffer, maxLength: bufferSize) guard bytesRead > 0 else { self.errorHandler(NetworkingError.InputError("Failed to read message bytes, got \(bytesRead)")) readState = .waitingForLength // 出错后重置状态 return } // 将读取的字节加入当前消息 currentMessage.append(buffer[0..<bytesRead]) remainingBytes -= bytesRead if remainingBytes == 0 { // 消息读取完成,转换为字符串 guard let message = String(data: currentMessage, encoding: .utf8) else { log.error(self.buildLogMessage("Failed to convert message data to UTF-8")) readState = .waitingForLength return } self.handleMessage(message) // 重置状态,等待下一条消息 readState = .waitingForLength } else { // 更新状态,若还有字节可用则继续读取 readState = .readingMessage(remainingBytes: remainingBytes, currentMessage: currentMessage) if istr.hasBytesAvailable { continueReadingMessage(remainingBytes: &remainingBytes, currentMessage: ¤tMessage) } } } private func cleanupInputStream() { self.inputStream?.close() self.inputStream = nil log.info(self.buildLogMessage("InputStream closed")) }
关键改进点:
- 状态管理:用
ReadState枚举明确当前读取阶段,避免重复读取长度字段; - 内存安全:用
Data累积消息内容,替代手动管理UnsafePointer的内存操作,减少崩溃风险; - 合理缓冲区:使用固定大小的缓冲区分批读取,避免因超大消息导致的内存压力;
- 线程安全:用串行队列保护流操作和状态变量,确保所有读取逻辑按顺序执行;
- 错误恢复:在读取出错或流结束时,及时重置状态并清理资源。
额外注意事项
- 确保
handleMessage方法线程安全:如果它需要更新UI或访问主线程资源,要通过DispatchQueue.main.async切换线程; - 处理流的错误和结束事件:避免因流异常导致的状态不一致;
- 测试边界情况:比如消息长度为0、超大消息、流中断等场景。
内容的提问来源于stack exchange,提问作者Simon Hessner
相关产品推荐
相关产品推荐

