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

Swift中MultipeerConnectivity大字节高效传输问题解决

问题描述

目标是在设备间生成并流式传输大数据(最大1.5GB)到磁盘,采用分块传输优化内存占用。当前基于MultipeerConnectivity的实现存在以下问题:

  • 发送端发送约17000字节后停止,接收端仅读取约2000字节后中断
  • 使用.current队列调度流时,StreamDelegate回调无法触发

当前核心实现

MultipeerService 核心逻辑

class MultipeerService: NSObject, ObservableObject, MCSessionDelegate, MCNearbyServiceBrowserDelegate, MCNearbyServiceAdvertiserDelegate {
    var inputStreamDelegate: DataReceiver?
    ...
    func startStreaming(withName streamName: String, to peer: MCPeerID) throws -> OutputStream {
        return try session.startStream(withName: streamName, toPeer: peer)
    }

    public func session(_ session: MCSession, didReceive stream: InputStream, withName streamName: String, fromPeer peerID: MCPeerID) {
        log.info("Started receiving stream '\(streamName)'")
        self.inputStreamDelegate?.readStream(stream, withName: streamName)
    }
}

DataSender 发送端实现

class DataSender: NSObject, StreamDelegate, ObservableObject {
    private var multipeerService: MultipeerService
    private var outputStream: OutputStream?
    private let chunkSize = 1024 // 1KB chunk size
    private var totalSize: Int = 0
    @Published var bytesSent: Int = 0

    init(ms: MultipeerService = MultipeerService.shared) {
        self.multipeerService = ms
        super.init()
    }

    // Called from UI
    func startStreaming(peer: MCPeerID, sizeMB: Int) {
        do {
            let outputStream = try multipeerService.startStreaming(withName: "Test", to: peer)
            self.outputStream = outputStream
            outputStream.delegate = self
            outputStream.schedule(in: .main, forMode: .common)
            outputStream.open()

            self.totalSize = sizeMB * 1024 * 1024 // convert MB to bytes
            self.bytesSent = 0
        } catch {
            fatalError("Error starting stream: \(error)")
        }
    }

    private func sendRandomData() {
        guard let outputStream = outputStream else { return }

        while outputStream.hasSpaceAvailable && bytesSent < totalSize {
            let currentChunkSize = min(chunkSize, totalSize - bytesSent)
            var buffer = [UInt8](repeating: 0, count: currentChunkSize)
            _ = SecRandomCopyBytes(kSecRandomDefault, currentChunkSize, &buffer)

            let bytesWritten = outputStream.write(buffer, maxLength: currentChunkSize)
            if bytesWritten > 0 {
                DispatchQueue.main.async {
                    self.bytesSent += bytesWritten
                }
                print("Bytes written: \(bytesWritten) | Total: \(bytesSent)")
            } else {
                if let error = outputStream.streamError {
                    print("Stream Error: \(error)")
                    closeStream()
                    return
                }
            }
        }

        if bytesSent >= totalSize {
            closeStream()
        }
    }

    private func closeStream() {
        print("Closing stream")
        outputStream?.close()
        outputStream?.remove(from: .main, forMode: .default)
        outputStream = nil
    }

    // StreamDelegate methods
    func stream(_ aStream: Stream, handle eventCode: Stream.Event) {
        print("Stream called: \(eventCode)")
        switch eventCode {
        case .hasSpaceAvailable:
            print("hasSpaceAvailable")
            sendRandomData()
        case .endEncountered:
            print("Stream end encountered")
            closeStream()
        case .errorOccurred:
            if let error = aStream.streamError {
                print("Stream Error: \(error)")
            }
            closeStream()
        default:
            break
        }
    }
}

DataReceiver 接收端实现

class DataReceiver: NSObject, StreamDelegate, ObservableObject {
    private var multipeerService: MultipeerService
    private var inputStream: InputStream?
    private var fileHandle: FileHandle?

    @Published var received: Int = 0

    init(ms: MultipeerService = MultipeerService.shared) {
        self.multipeerService = ms
        super.init()

        self.multipeerService.inputStreamDelegate = self
    }

    // Called by MultipeerService when stream is received
    public func readStream(_ stream: InputStream, withName streamName: String) {
        stream.delegate = self
        stream.schedule(in: .main, forMode: .default)
        stream.open()
        self.inputStream = stream

        // Prepare to write to file
        if let fileURL = createFileInLibrary(withName: streamName) {
            do {
                fileHandle = try FileHandle(forWritingTo: fileURL)
            } catch {
                print("Error opening file handle: \(error)")
            }
        }
    }

    // StreamDelegate method to handle input stream events
    func stream(_ aStream: Stream, handle eventCode: Stream.Event) {
        print("Stream function called. EventCode: \(eventCode)")
        switch eventCode {
        case .hasBytesAvailable:
            if aStream == inputStream {
                readAvailableBytes(stream: aStream as! InputStream)
            }
        case .endEncountered:
            closeStream()
        case .errorOccurred:
            if let error = aStream.streamError {
                print("Stream Error: \(error)")
            }
            closeStream()
        default:
            break
        }
    }

    private func readAvailableBytes(stream: InputStream) {
        let bufferSize = 1024
        var buffer = [UInt8](repeating: 0, count: bufferSize)

        while stream.hasBytesAvailable {
            let numberOfBytesRead = stream.read(&buffer, maxLength: bufferSize)
            print("Number of bytes read: \(numberOfBytesRead)")
            DispatchQueue.main.async {
                self.received += numberOfBytesRead
            }

            if numberOfBytesRead < 0 {
                if let error = stream.streamError {
                    print("Stream read error: \(error)")
                    return
                }
            }

            // Write to file
            if let fileHandle = fileHandle {
                fileHandle.write(Data(bytes: buffer, count: numberOfBytesRead))
            }
        }
    }

    private func closeStream() {
        inputStream?.close()
        inputStream?.remove(from: .current, forMode: .default)
        inputStream = nil

        fileHandle?.closeFile()
        fileHandle = nil
    }

    private func createFileInLibrary(withName name: String) -> URL? {
        let fileManager = FileManager.default
        guard let libraryDirectory = fileManager.urls(for: .libraryDirectory, in: .userDomainMask).first else {
            print("Could not find library directory")
            return nil
        }

        let fileURL = libraryDirectory.appendingPathComponent("\(name).data")
        fileManager.createFile(atPath: fileURL.path, contents: nil, attributes: nil)

        return fileURL
    }
}

错误现象

接收端控制台输出:

Stream function called. EventCode: NSStreamEvent(rawValue: 1)
Stream function called. EventCode: NSStreamEvent(rawValue: 2)
Number of bytes read: 1024
Number of bytes read: 1024
Number of bytes read: 142

解决方案与优化建议

一、修复传输中断问题

1. 发送端避免阻塞Runloop

当前sendRandomData中的while循环会持续占用主线程Runloop,导致后续Stream.Event无法触发。修改为单次写入后退出,等待下一次.hasSpaceAvailable事件或异步触发下一次写入:

private func sendRandomData() {
    guard let outputStream = outputStream, bytesSent < totalSize else { return }
    
    if !outputStream.hasSpaceAvailable { return }
    
    let currentChunkSize = min(chunkSize, totalSize - bytesSent)
    var buffer = [UInt8](repeating: 0, count: currentChunkSize)
    _ = SecRandomCopyBytes(kSecRandomDefault, currentChunkSize, &buffer)
    
    let bytesWritten = outputStream.write(buffer, maxLength: currentChunkSize)
    if bytesWritten > 0 {
        DispatchQueue.main.async {
            self.bytesSent += bytesWritten
        }
        print("Bytes written: \(bytesWritten) | Total: \(bytesSent)")
        // 还有数据未发送时,异步触发下一次写入,避免阻塞Runloop
        if bytesSent < totalSize {
            DispatchQueue.main.async { [weak self] in
                self?.sendRandomData()
            }
        } else {
            closeStream()
        }
    } else {
        if let error = outputStream.streamError {
            print("Stream Error: \(error)")
            closeStream()
        }
    }
}

2. 接收端处理读取边界情况

补充stream.read返回0(流结束)的处理逻辑,同时校验读取字节数的有效性:

private func readAvailableBytes(stream: InputStream) {
    let bufferSize = 1024
    var buffer = [UInt8](repeating: 0, count: bufferSize)

    while stream.hasBytesAvailable {
        let numberOfBytesRead = stream.read(&buffer, maxLength: bufferSize)
        
        if numberOfBytesRead == 0 {
            // 流已结束
            closeStream()
            return
        }
        
        if numberOfBytesRead < 0 {
            if let error = stream.streamError {
                print("Stream read error: \(error)")
            }
            closeStream()
            return
        }
        
        print("Number of bytes read: \(numberOfBytesRead)")
        DispatchQueue.main.async {
            self.received += numberOfBytesRead
        }

        // Write to file
        if let fileHandle = fileHandle {
            fileHandle.write(Data(bytes: buffer, count: numberOfBytesRead))
        }
    }
}

3. 发送端主动触发首次写入

在startStreaming中打开流后手动调用一次sendRandomData,避免等待事件的延迟:

func startStreaming(peer: MCPeerID, sizeMB: Int) {
    do {
        let outputStream = try multipeerService.startStreaming(withName: "Test", to: peer)
        self.outputStream = outputStream
        outputStream.delegate = self
        outputStream.schedule(in: .main, forMode: .common)
        outputStream.open()

        self.totalSize = sizeMB * 1024 * 1024
        self.bytesSent = 0
        // 主动触发首次数据发送
        sendRandomData()
    } catch {
        fatalError("Error starting stream: \(error)")
    }
}

二、修复StreamDelegate未触发问题(非主线程调度)

.current队列默认没有活跃的Runloop,无法触发代理回调。需创建自定义串行队列并绑定Runloop:

发送端修改调度逻辑

// DataSender中添加属性
private let streamQueue = DispatchQueue(label: "com.example.stream.send")
private var streamRunloop: RunLoop?

func startStreaming(peer: MCPeerID, sizeMB: Int) {
    do {
        let outputStream = try multipeerService.startStreaming(withName: "Test", to: peer)
        self.outputStream = outputStream
        outputStream.delegate = self
        
        // 在自定义队列中调度流并启动Runloop
        streamQueue.async {
            self.streamRunloop = RunLoop.current
            outputStream.schedule(in: self.streamRunloop!, forMode: .default)
            outputStream.open()
            
            // 定期唤醒Runloop,避免永久阻塞
            while self.outputStream != nil && self.streamRunloop != nil {
                self.streamRunloop!.run(until: Date().addingTimeInterval(0.1))
            }
        }

        self.totalSize = sizeMB * 1024 * 1024
        self.bytesSent = 0
    } catch {
        fatalError("Error starting stream: \(error)")
    }
}

// 修改closeStream方法,终止Runloop
private func closeStream() {
    print("Closing stream")
    streamQueue.async {
        self.outputStream?.close()
        self.outputStream?.remove(from: self.streamRunloop!, forMode: .default)
        self.outputStream = nil
        self.streamRunloop = nil
    }
}

接收端同理修改调度逻辑

// DataReceiver中添加属性
private let streamQueue = DispatchQueue(label: "com.example.stream.receive")
private var streamRunloop: RunLoop?

public func readStream(_ stream: InputStream, withName streamName: String) {
    stream.delegate = self
    
    // 在自定义队列中调度流并启动Runloop
    streamQueue.async {
        self.streamRunloop = RunLoop.current
        stream.schedule(in: self.streamRunloop!, forMode: .default)
        stream.open()
        self.inputStream = stream
        
        while self.inputStream != nil && self.streamRunloop != nil {
            self.streamRunloop!.run(until: Date().addingTimeInterval(0.1))
        }
    }

    // 初始化文件句柄
    if let fileURL = createFileInLibrary(withName: streamName) {
        do {
            fileHandle = try FileHandle(forWritingTo: fileURL)
        } catch {
            print("Error opening file handle: \(error)")
        }
    }
}

// 修改closeStream方法
private func closeStream() {
    streamQueue.async {
        self.inputStream?.close()
        self.inputStream?.remove(from: self.streamRunloop!, forMode: .default)
        self.inputStream = nil
        self.streamRunloop = nil
    }
    
    fileHandle?.closeFile()
    fileHandle = nil
}

三、性能优化建议

  1. 增大分块大小:将chunkSize从1KB提升至64KB或128KB,减少IO操作次数,提升传输效率。
  2. 后台传输支持:在Info.plist中添加NSBonjourServices和UIBackgroundModes(包含processing和networking),允许应用在后台进行传输。
  3. 错误重试机制:针对传输中的可恢复错误,实现有限次数的重试逻辑,避免单次错误终止整个传输。
  4. 进度回调优化:降低bytesSent和received的更新频率(例如每1MB更新一次),减少主线程压力。
  5. 异步文件写入:将文件写入操作放到后台队列,避免阻塞流读取的Runloop:
// 接收端修改文件写入部分
if let fileHandle = fileHandle {
    DispatchQueue.global(qos: .background).async {
        fileHandle.write(Data(bytes: buffer, count: numberOfBytesRead))
    }
}
  1. 流状态校验:在打开流、写入/读取前后增加流状态校验,提前发现异常并处理。

内容的提问来源于stack exchange,提问作者Jan Blažek

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 11:52:03