基于Combine实现CoreBluetooth扫描与订阅数联动启停
用Combine封装CoreBluetooth:实现订阅驱动的扫描启停
你遇到的核心问题是多订阅场景下的扫描生命周期管理——handleEvents会为每个新订阅触发receiveSubscription,若不共享上游会导致重复调用扫描启动;同时需要准确跟踪订阅数量,确保仅在最后一个订阅取消时停止扫描。以下是两种可落地的解决方案:
方案一:用share() + 订阅计数器(轻量改造)
无需自定义Publisher,通过共享上游和线程安全计数器实现逻辑,适合快速适配现有代码。
步骤1:添加订阅计数管理
在Bluetooth类中新增线程安全的计数变量与保护队列:
final class Bluetooth { // ... 原有代码 ... private var subscriberCount = 0 private let subscriberQueue = DispatchQueue(label: "com.mycompany.bluetooth.subscriberQueue") // ... 原有代码 ... }
步骤2:改造scan()方法
整合蓝牙状态校验、设备发现流,并添加订阅触发的扫描逻辑:
func scan() -> AnyPublisher<[UUID: DiscoverResult], Bluetooth.Error> { // 蓝牙状态流转:输出就绪状态或错误 let isBluetoothReady = delegateWrapper.state .map { state -> Result<Bool, Bluetooth.Error> in switch state { case .poweredOn: return .success(true) case .unknown, .resetting: return .success(false) case .poweredOff: return .failure(.poweredOff) case .unsupported: return .failure(.unsupported) case .unauthorized: return .failure(.unauthorized(CBCentralManager.authorization)) @unknown default: return .success(false) } } .share() // 设备发现流:仅在蓝牙就绪时输出数据 let discoveryStream = isBluetoothReady .compactMap { result in switch result { case .success(true): return delegateWrapper.didDiscoverPeripheral.eraseToAnyPublisher() default: return Empty().eraseToAnyPublisher() } } .switchToLatest() .setFailureType(to: Bluetooth.Error.self) // 合并错误流与发现流 let combinedStream = isBluetoothReady .compactMap { result -> AnyPublisher<[UUID: DiscoverResult], Bluetooth.Error>? in switch result { case .failure(let error): return Fail(error: error).eraseToAnyPublisher() default: return nil } } .switchToLatest() .merge(with: discoveryStream) // 添加订阅启停逻辑 return combinedStream .handleEvents( receiveSubscription: { [weak self] _ in guard let self = self else { return } self.subscriberQueue.sync { self.subscriberCount += 1 if self.subscriberCount == 1 { // 第一个订阅:启动扫描(二次校验蓝牙状态) self.subscriberQueue.async { [weak self] in guard let self = self, self.central.state == .poweredOn else { return } self.central.scanForPeripherals(withServices: nil, options: nil) } } } }, receiveCancel: { [weak self] in guard let self = self else { return } self.subscriberQueue.sync { self.subscriberCount -= 1 if self.subscriberCount == 0 { // 最后一个订阅取消:停止扫描并清空缓存 self.central.stopScan() self.delegateWrapper.peripherals.removeAll() self.delegateWrapper.didDiscoverPeripheral.send([:]) } } } ) .share() // 关键:让所有订阅共享同一上游,避免计数混乱 .eraseToAnyPublisher() }
核心要点
share()的作用:确保所有订阅者复用同一个流实例,避免每个订阅都触发独立的扫描启动。- 线程安全计数:用串行队列保护
subscriberCount,规避多线程竞态问题。 - 二次状态校验:启动扫描前再次确认蓝牙状态,防止订阅时蓝牙尚未就绪的情况。
方案二:自定义Publisher(完全可控)
若需要更精细的生命周期控制(如需求调度、错误精准传播),可自定义Publisher直接管理订阅与扫描状态。
自定义扫描Publisher与Subscription
private class ScanPublisher: Publisher { typealias Output = [UUID: DiscoverResult] typealias Failure = Bluetooth.Error private weak var bluetooth: Bluetooth? private let discoveryStream: AnyPublisher<Output, Failure> private let stateStream: AnyPublisher<CBManagerState, Never> private var subscriberCount = 0 private let queue = DispatchQueue(label: "com.mycompany.bluetooth.scanPublisher") init(bluetooth: Bluetooth) { self.bluetooth = bluetooth self.discoveryStream = bluetooth.delegateWrapper.didDiscoverPeripheral .setFailureType(to: Failure.self) .eraseToAnyPublisher() self.stateStream = bluetooth.delegateWrapper.state.eraseToAnyPublisher() } func receive<S>(subscriber: S) where S : Subscriber, Failure == S.Failure, Output == S.Input { let subscription = ScanSubscription( subscriber: subscriber, publisher: self, discoveryStream: discoveryStream, stateStream: stateStream ) subscriber.receive(subscription: subscription) queue.sync { subscriberCount += 1 if subscriberCount == 1 { queue.async { [weak self] in guard let self = self, let bluetooth = self.bluetooth else { return } if bluetooth.central.state == .poweredOn { bluetooth.central.scanForPeripherals(withServices: nil, options: nil) } } } } } func cancelSubscription() { queue.sync { subscriberCount -= 1 if subscriberCount == 0 { bluetooth?.central.stopScan() bluetooth?.delegateWrapper.peripherals.removeAll() bluetooth?.delegateWrapper.didDiscoverPeripheral.send([:]) } } } } private class ScanSubscription<S: Subscriber>: Subscription where S.Input == [UUID: DiscoverResult], S.Failure == Bluetooth.Error { private var subscriber: S? private var discoveryCancellable: AnyCancellable? private var stateCancellable: AnyCancellable? private weak var publisher: ScanPublisher? init(subscriber: S, publisher: ScanPublisher, discoveryStream: AnyPublisher<[UUID: DiscoverResult], Bluetooth.Error>, stateStream: AnyPublisher<CBManagerState, Never>) { self.subscriber = subscriber self.publisher = publisher // 监听设备发现流 discoveryCancellable = discoveryStream .sink( receiveCompletion: { [weak self] completion in self?.subscriber?.receive(completion: completion) self?.cleanup() }, receiveValue: { [weak self] value in _ = self?.subscriber?.receive(value) } ) // 监听蓝牙状态错误 stateCancellable = stateStream .map { state -> Bluetooth.Error? in switch state { case .poweredOff: return .poweredOff case .unsupported: return .unsupported case .unauthorized: return .unauthorized(CBCentralManager.authorization) default: return nil } } .compactMap { $0 } .sink { [weak self] error in self?.subscriber?.receive(completion: .failure(error)) self?.cleanup() } } func request(_ demand: Subscribers.Demand) {} func cancel() { publisher?.cancelSubscription() cleanup() } private func cleanup() { subscriber = nil discoveryCancellable?.cancel() stateCancellable?.cancel() } }
在Bluetooth中使用自定义Publisher
func scan() -> AnyPublisher<[UUID: DiscoverResult], Bluetooth.Error> { ScanPublisher(bluetooth: self) .eraseToAnyPublisher() }
核心要点
- 完全可控的生命周期:直接管理订阅的添加与取消,精准控制扫描启停时机。
- 统一错误处理:在Subscription中整合蓝牙状态错误与设备发现流错误,逻辑更清晰。
- 无共享依赖:无需依赖
share(),每个订阅由Publisher独立跟踪,避免潜在的流共享问题。
额外注意事项
- 扫描参数定制:可在
scanForPeripherals中添加服务UUID过滤或扫描选项(如允许重复扫描),适配业务需求。 - 状态动态响应:可扩展状态流监听,在订阅存在时若蓝牙从
poweredOn变为其他状态,自动停止扫描,待状态恢复后重新启动。
内容的提问来源于stack exchange,提问作者corujautx
相关产品推荐
相关产品推荐

