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

并行执行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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 02:48:01