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

合并两个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()
        }
    }
}

额外需求(更新)

  1. 同步场景参考:同步场景下的实现方案如下:
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
  1. 优先级队列需求: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))

方案说明

  1. 优先级队列:基于最小堆实现,新值入队时自动完成排序,确保能快速获取当前最小的未处理值。
  2. 并发流消费:两个独立任务分别处理偶数和奇数流,互不阻塞,保证流的推送效率。
  3. 有序处理逻辑:循环检查队列头部是否匹配当前预期值,一旦匹配就处理并更新预期值,实现连续有序的处理。
  4. 性能优化:加入短暂休眠避免空轮询,也可以通过AsyncSemaphore或Continuation实现更高效的唤醒机制(有新值入队时立即唤醒处理任务)。

内容的提问来源于stack exchange,提问作者Luca Angeletti

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 16:06:06