并行执行Combine Publisher出现竞态条件,HealthKit查询结果匹配错误
问题核心原因
竞态问题的根源是Publishers.MergeMany的行为特性:它会将多个输入发布者的输出按结果返回的先后顺序下发,而非按照你传入的workoutSegments的原始顺序下发,所以最终和workoutSegments.publisher对齐时自然会出现数据错位。
修复方案
不要全局合并所有心率、步数的查询请求,改为针对每个独立的运动分段,将其对应的心率、步数查询和分段本身绑定,最后按原始顺序合并所有分段的结果即可保证顺序匹配:
// 提取Apple Watch自动生成的运动分段 let workoutSegments = (workout.workoutEvents ?? []).filter({ $0.type == .segment }) // 对每个分段单独发起并行的心率、步数查询,和分段本身绑定 let segmentPublishers = workoutSegments.map { segment in let interval = segment.dateInterval // 当前分段的心率查询 let hrPublisher = healthStore.statistic( for: HKQuantityType.quantityType(forIdentifier: .heartRate)!, with: .discreteAverage, from: interval.start, to: interval.end ) .assertNoFailure() // 当前分段的步数查询 let stepPublisher = healthStore.statistic( for: HKQuantityType.quantityType(forIdentifier: .stepCount)!, with: .cumulativeSum, from: interval.start, to: interval.end ) .assertNoFailure() // 单个分段的三个数据绑定,不会和其他分段混淆 return Publishers.Zip3( Just(segment).setFailureType(to: Never.self), stepPublisher, hrPublisher ) } // 严格按原始分段顺序合并所有结果 Publishers.ConcatenateMany(segmentPublishers) .receive(on: DispatchQueue.main) .sink(receiveValue: { segment, steps, hrs in let d = SGWorkoutEvent( type: segment.type, splitActiveDurationQuantity: segment.splitDuration, splitDistanceQuantity: segment.splitDistance, totalDistanceQuantity: segment.totalDistanceQuantity, splitMeasuringSystem: segment.splitMeasuringSystem, steps: steps.sumQuantity(), heartRate: hrs.averageQuantity() ) self.paces.append(d) }) .store(in: &bag)
方案说明
- 每个运动分段的心率、步数查询和当前分段强绑定,不管两个查询谁先返回,
Zip3都会等两者都完成后和当前分段配对,不会出现跨分段的数据错位 Publishers.ConcatenateMany会严格按照workoutSegments的原始顺序下发结果,只有前一个分段的所有查询完成后才会处理下一个分段的结果,完全保证顺序匹配。
额外优化点
你当前的HKHealthStore.statistic扩展封装存在小缺陷:查询报错时仅打印了日志,没有向发布者下发错误事件,会导致出错时发布者永远不会终止,建议补充错误下发逻辑:
public func statistic(for type: HKQuantityType, with options: HKStatisticsOptions, from startDate: Date, to endDate: Date, _ limit: Int = HKObjectQueryNoLimit) -> AnyPublisher<HKStatistics, Error> { let subject = PassthroughSubject<HKStatistics, Error>() let predicate = HKStatisticsQuery.predicateForSamples(withStart: startDate, end: endDate, options: [.strictEndDate, .strictStartDate]) let query = HKStatisticsQuery(quantityType: type, quantitySamplePredicate: predicate, options: options, completionHandler: { (query, statistics, error) in if let error = error { hkCombineLogger.error("Error fetching statistics \(error.localizedDescription)") subject.send(completion: .failure(error)) return } subject.send(statistics!) subject.send(completion: .finished) }) self.execute(query) return subject.eraseToAnyPublisher() }
内容的提问来源于stack exchange,提问作者user4296481
相关产品推荐
相关产品推荐

