在Swift并发模型里,异步序列让我们可以像处理同步集合一样逐条消费异步产生的数据。但当数据源本身可能失败时,普通去重逻辑往往束手无策。AsyncThrowingDistinctSequence作为AsyncSequence的一个变体,专门解决边去重边抛错的需求,它内部维护已见元素的集合,并在每次next()调用时允许抛出错误,从而把网络抖动、解码异常等故障暴露给上层而不是默默丢弃。

AsyncSequence与去重算子的底层机制
AsyncSequence协议要求类型提供makeAsyncIterator方法,返回一个遵守AsyncIteratorProtocol的对象,其核心是异步的next()函数。标准库通过扩展提供了不少算子,例如map、filter、compactMap,但去重并不是所有变体都内置。对于会抛错的场景,编译器需要我们明确使用throwing版本,否则错误类型会被推断为Never,导致无法在去重过程中向上抛异常。
AsyncThrowingDistinctSequence的实现思路并不复杂:迭代器持有一个Set用来记录已发射的元素哈希值,每次从上游拉取一个元素,先尝试计算哈希并查询集合。如果已存在就跳过,不存在则插入并产出。关键在于上游的next()可能是抛错的,因此它自己的next()也标记为throws,调用方必须用try await来消费。这种结构让去重状态机和错误通道解耦,我们不必在业务代码里写do-catch包裹整个循环。
对比自己手写去重,手动方案通常要在for await里声明一个var seen = Set(),然后每次判断。一旦上游是AsyncThrowingStream,你就得在循环体里写try,一旦忘记处理某个调用点,编译器虽能报错,但逻辑散落各处。使用AsyncThrowingDistinctSequence可以把这块通用逻辑收敛到标准库,既减少bug也提升可读性。
用AsyncThrowingDistinctSequence处理实时日志去重与错误
假设我们有一个异步日志源,它从多个传感器读取事件,某些事件因为校验失败需要抛出。我们希望对相同ID的事件去重,同时把校验错误传递给订阅者而不是中断整个流。下面代码演示如何构造并使用该序列。
import Foundation
struct LogEvent: Hashable {
let id: String
let value: Int
}
func sensorStream() -> AsyncThrowingStream<LogEvent, Error> {
AsyncThrowingStream { continuation in
let items = [
LogEvent(id: "a", value: 1),
LogEvent(id: "a", value: 1),
LogEvent(id: "b", value: 2)
]
Task {
for item in items {
if item.value < 0 {
continuation.finish(throwing: NSError(domain: "bad", code: 1))
return
}
continuation.yield(item)
}
// 模拟一个错误
continuation.finish(throwing: NSError(domain: "done", code: 0))
}
}
}
func demo() async {
let distinct = sensorStream().distinct()
do {
for try await event in distinct {
print("emit: (event.id)")
}
} catch {
print("stream ended with error: (error)")
}
}
上面示例中,distinct()是标准库为AsyncThrowingStream提供的去重算子,返回类型即AsyncThrowingDistinctSequence。我们循环里用try await消费,重复ID的第二个a被自动过滤,而流结束时抛出的错误被catch捕获并打印。这样去重与错误处理各司其职,调用方逻辑非常清晰。
需要注意,distinct()默认使用元素自身的Hashable实现。如果事件结构体很大或哈希成本高,可以自定义哈希方式,或者先map成轻量键再distinct。另外,该序列是惰性的,只有消费端请求next时才会从上游拉数据,因此不会因为缓冲导致内存膨胀,这对长周期日志流尤其重要。
性能与边界情况分析
在高频数据场景下,AsyncThrowingDistinctSequence的Set会带来O(1)平均查找开销,但内存随去重跨度增长。如果业务只需要最近N个去重,可以改用滑动窗口自己封装,或者定期清理Set。由于它是串行迭代,不会存在多线程同时写集合的问题,这比手动在并发任务里加锁简单得多。
另一个边界是错误发生位置。如果错误在上游yield之前抛出,distinct序列会立刻把错误向上传递,已发出的元素不受影响;如果错误发生在yield之后、next返回前,同样遵循AsyncThrowingStream的契约。开发者应当把可能失败的解码操作放在上游流里,而不是在去重之后,否则重复元素可能已经被发出,错误却迟到,造成状态不一致。
从架构角度看,把AsyncThrowingDistinctSequence放在管道前部,紧挨数据源,能保证尽早过滤噪声并暴露故障;若放在多个map之后,则要去重的是变换后的模型,更符合业务语义。两种摆法没有绝对优劣,取决于你对“重复”的定义落在哪一层。掌握这套类型后,Swift异步数据清洗会变得更声明式且安全。
AsyncSequenceAsyncThrowingDistinctSequenceSwift去重修改时间:2026-08-14 22:24:28