合并两个AsyncThrowingStream并有序处理的实现方案问询
问题描述
背景
我有两个AsyncThrowingStream,分别推送前5个非负偶数和奇数:
var evens = stream(start: 0) // 0, 2, 4, 6, 8 var odds = stream(start: 1) // 1, 3, 5, 7, 9
每个值会在1-3秒后推送。
并发需求
由于两个流推送下一个值需要时间,我希望它们并发运行,示例代码如下:
Task { let evens = stream(start: 0) for try await even in evens { } } Task { let odds = stream(start: 1) for try await odd in odds { } }
有序处理需求
我有如下处理函数:
func process(value: Int) { print(value) }
需要按顺序调用process(value:)处理所有生成的整数,预期调用顺序:
process(value: 0) process(value: 1) process(value: 2) process(value: 3) process(value: 4) process(value: 5) process(value: 6) process(value: 7) process(value: 8) process(value: 9)
补充测试代码
为方便测试,提供生成两个流的模拟实现:
func stream(start: Int) -> AsyncThrowingStream<Int, Error> { .init { continuation in Task { for index in stride(from: start, to: 10, by: 2) { continuation.yield(index) let random = UInt64.random(in: (1...3)) try await Task.sleep(nanoseconds: random * 1_000_000_000) } continuation.finish() } } }
额外需求(更新)
- 同步场景参考:同步场景下的实现方案如下:
let evens = [0, 2, 4, 6, 8] let odds = [1, 3, 5, 7, 9] zip(evens, odds) .map { [$0, $1] } .flatMap { $0 } .forEach(process) 输出: 0 1 2 3 4 5 6 7 8 9
- 优先级队列需求:
process(value:)必须在最后一个已处理整数的后继值被任意一个流生成后立即调用:- 调用参数必须遵循非负整数的自然顺序(0,1,2,3,4...);
- 必须尽可能早地调用(一旦符合条件的下一个整数生成)。
示例流程:
evens推送0 → 调用process(value:0) odds推送1 → 调用process(value:1) odds推送3 → 3被暂存 evens推送2 → 调用process(value:2),接着调用process(value:3)
解决方案
可以通过**优先级队列(最小堆)**结合Swift并发任务调度实现,核心思路:
- 并发消费两个流,将收到的值存入优先级队列;
- 单独启动任务,持续检查队列头部是否为当前需要处理的下一个值,满足条件则取出处理,直到所有流完成且队列清空。
实现代码
import Foundation func stream(start: Int) -> AsyncThrowingStream<Int, Error> { .init { continuation in Task { for index in stride(from: start, to: 10, by: 2) { continuation.yield(index) let random = UInt64.random(in: (1...3)) try await Task.sleep(nanoseconds: random * 1_000_000_000) } continuation.finish() } } } func process(value: Int) { print(value) } func mergeAndProcessOrdered() async throws { // 最小堆实现的优先级队列,保证每次取出当前最小的未处理值 var priorityQueue = PriorityQueue<Int>(order: <) // 记录当前预期要处理的下一个值 var nextExpected = 0 // 标记两个流是否已完成推送 var evensCompleted = false var oddsCompleted = false // 并发消费偶数流 let evensTask = Task { for try await value in stream(start: 0) { priorityQueue.enqueue(value) } evensCompleted = true } // 并发消费奇数流 let oddsTask = Task { for try await value in stream(start: 1) { priorityQueue.enqueue(value) } oddsCompleted = true } // 循环检查并处理符合顺序要求的值 while !evensCompleted || !oddsCompleted || !priorityQueue.isEmpty { // 连续处理所有符合当前预期的连续值 while let top = priorityQueue.peek(), top == nextExpected { let value = priorityQueue.dequeue()! process(value: value) nextExpected += 1 } // 短暂休眠避免空轮询占用CPU资源 try await Task.sleep(nanoseconds: 100_000_000) } // 等待两个消费任务完成 try await evensTask.value try await oddsTask.value } // 简易最小堆优先级队列实现 struct PriorityQueue<Element: Comparable> { private var elements: [Element] private let order: (Element, Element) -> Bool init(order: @escaping (Element, Element) -> Bool) { self.elements = [] self.order = order } mutating func enqueue(_ element: Element) { elements.append(element) bubbleUp(from: elements.count - 1) } mutating func dequeue() -> Element? { guard !elements.isEmpty else { return nil } if elements.count == 1 { return elements.removeLast() } let front = elements[0] elements[0] = elements.removeLast() bubbleDown(from: 0) return front } func peek() -> Element? { elements.first } var isEmpty: Bool { elements.isEmpty } private mutating func bubbleUp(from index: Int) { var childIndex = index let child = elements[childIndex] var parentIndex = (childIndex - 1) / 2 while childIndex > 0 && order(child, elements[parentIndex]) { elements[childIndex] = elements[parentIndex] childIndex = parentIndex parentIndex = (childIndex - 1) / 2 } elements[childIndex] = child } private mutating func bubbleDown(from index: Int) { let count = elements.count var parentIndex = index while true { let leftChildIndex = 2 * parentIndex + 1 let rightChildIndex = 2 * parentIndex + 2 var candidateIndex = parentIndex if leftChildIndex < count && order(elements[leftChildIndex], elements[candidateIndex]) { candidateIndex = leftChildIndex } if rightChildIndex < count && order(elements[rightChildIndex], elements[candidateIndex]) { candidateIndex = rightChildIndex } if candidateIndex == parentIndex { break } elements.swapAt(parentIndex, candidateIndex) parentIndex = candidateIndex } } } // 测试调用 Task { try await mergeAndProcessOrdered() } // 命令行环境下保持程序运行 RunLoop.main.run(until: Date().addingTimeInterval(30))
方案说明
- 优先级队列:基于最小堆实现,新值入队时自动完成排序,确保能快速获取当前最小的未处理值。
- 并发流消费:两个独立任务分别处理偶数和奇数流,互不阻塞,保证流的推送效率。
- 有序处理逻辑:循环检查队列头部是否匹配当前预期值,一旦匹配就处理并更新预期值,实现连续有序的处理。
- 性能优化:加入短暂休眠避免空轮询,也可以通过
AsyncSemaphore或Continuation实现更高效的唤醒机制(有新值入队时立即唤醒处理任务)。
内容的提问来源于stack exchange,提问作者Luca Angeletti
相关产品推荐
相关产品推荐

