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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 20:35:05