Documentation
¶
Overview ¶
Package async 提供一次性结果、完成信号、连续流和结构化异步作用域。
Package async 将不同异步语义拆分为独立类型:
- Promise/Future:非泛型、一次性、可重放的 Result;
- Completer/Signal:不携带 Result 的生命周期完成通知;
- Emitter/Stream:连续 Result 的单消费流;
- Scope/Spawn:绑定宿主生命周期的后台任务取消、汇合与统计;
- Race、FirstSuccess、All、AllSettled、Zip2、Map、FlatMap 和 Timeout: 基于完成订阅的 Future 组合器。
Future 内部保存完成结果,并通过 OnComplete 在完成者 goroutine 中直接通知订阅者; 多个 Wait、TryGet 或 OnComplete 消费者读取的是同一个可重放结果,不会竞争消费, 也不会为每个订阅者启动等待 goroutine。
Stream 采用单消费语义;多个消费者会竞争元素。需要广播时应使用 event 或上层消息 设施。Scope.Close 只负责取消 Context 和等待协作式退出,不能强制终止 goroutine。
根包 core 的 Submit、Post、Spawn 和 ContinueOn 建立在这些能力之上。
Index ¶
- Variables
- func NewPromise(completionExecutorID ...ExecutorID) (Promise, Future)
- func NewSignal() (Completer, Signal)
- func NewStream(buffer ...int) (Emitter, Stream)
- type Completer
- type Emitter
- type ExecutorID
- type Future
- func All(futures ...Future) Future
- func AllSettled(futures ...Future) Future
- func FirstSuccess(futures ...Future) Future
- func FlatMap(future Future, fn func(Result) Future) Future
- func Map(future Future, fn func(Result) Result) Future
- func Race(futures ...Future) Future
- func Rejected(err error) Future
- func Resolved(ret Result) Future
- func Spawn(scope *Scope, task func(context.Context) Result) Future
- func SpawnVoid(scope *Scope, task func(context.Context)) Future
- func Timeout(ctx context.Context, future Future, duration time.Duration) Future
- func Zip2(first, second Future) Future
- func (future Future) CompletionExecutorID() ExecutorID
- func (future Future) Context(ctx context.Context) context.Context
- func (future Future) Done() <-chan struct{}
- func (future Future) ID() uint64
- func (future Future) IsNil() bool
- func (future Future) OnComplete(callback func(Result)) Subscription
- func (future Future) TryGet() (ret Result, ok bool)
- func (future Future) Wait(ctx context.Context) Result
- type Pair
- type Promise
- type Result
- type Scope
- type ScopeStats
- type Signal
- type Stream
- type Subscription
- type WaitGuard
Constants ¶
This section is empty.
Variables ¶
var ( ErrAsync = fmt.Errorf("%w: async", exception.ErrCore) // ErrAsync 是异步模块错误的共同根错误。 ErrScopeClosed = fmt.Errorf("%w: scope closed", ErrAsync) // ErrScopeClosed 表示异步作用域已经关闭。 ErrNoCandidates = fmt.Errorf("%w: no future candidates", ErrAsync) // ErrNoCandidates 表示组合器没有可用的候选 Future。 ErrNoFutureSucceeded = fmt.Errorf("%w: no future succeeded", ErrAsync) // ErrNoFutureSucceeded 表示所有候选 Future 均失败。 ErrFutureTimeout = fmt.Errorf("%w: future timeout", ErrAsync) // ErrFutureTimeout 表示 Future 等待超时。 )
Functions ¶
func NewPromise ¶ added in v0.4.27
func NewPromise(completionExecutorID ...ExecutorID) (Promise, Future)
NewPromise 创建一对生产者 Promise 和消费者 Future。
completionExecutorID 可选;仅使用第一个值。它应表示完成该 Future 必须运行的 Runtime 执行器,外部 I/O、后台 goroutine 或未知来源使用零值。
Types ¶
type Completer ¶ added in v0.4.27
type Completer struct {
// contains filtered or unexported fields
}
Completer 是无结果完成信号的生产者端。
type Emitter ¶ added in v0.4.27
type Emitter struct {
// contains filtered or unexported fields
}
Emitter 是 Stream 的生产者端。通常只应由一个 goroutine 使用。
type ExecutorID ¶ added in v0.4.27
type ExecutorID uint64
ExecutorID 标识进程内的异步结果完成执行器。零值表示结果由外部执行器完成或归属未知。
ExecutorID 只用于运行时阻塞检测,不是持久化 ID,也不应跨进程传输。
func GenExecutorID ¶ added in v0.4.27
func GenExecutorID() ExecutorID
GenExecutorID 生成进程内唯一的非零异步执行器 ID。
type Future ¶ added in v0.4.27
type Future struct {
// contains filtered or unexported fields
}
Future 是一次性异步结果的只读消费者视图。
完成结果会保存在共享状态中,Wait、TryGet 和 OnComplete 均具有重放语义;多个 消费者不会竞争同一个结果。
func AllSettled ¶ added in v0.4.27
AllSettled 按输入顺序返回全部 Result,无论各项成功或失败。 空输入成功返回空 []Result。
func FirstSuccess ¶ added in v0.4.27
FirstSuccess 返回第一个成功结果;所有候选均失败时返回包装 ErrNoFutureSucceeded 的错误。
func Race ¶ added in v0.4.27
Race 返回第一个完成的 Future 结果。候选为空时返回 ErrNoCandidates。
Race 只解除失败者订阅,不取消可能被其他调用者共享的生产任务。
func Timeout ¶ added in v0.4.27
Timeout 返回在源 Future、ctx 取消或 duration 到期三者中最先完成的结果。 它只解除源 Future 的订阅,不取消源任务。
func (Future) CompletionExecutorID ¶ added in v0.4.27
func (future Future) CompletionExecutorID() ExecutorID
CompletionExecutorID 返回完成 Future 所依赖的执行器 ID。
func (Future) Done ¶ added in v0.4.27
func (future Future) Done() <-chan struct{}
Done 返回 Future 完成时关闭的共享频道;零值 Future 会导致 panic。
func (Future) OnComplete ¶ added in v0.4.27
func (future Future) OnComplete(callback func(Result)) Subscription
OnComplete 订阅 Future 完成。Future 已完成时 callback 会在调用者 goroutine 中立即执行。 callback 为 nil 时会导致 panic。
type Promise ¶ added in v0.4.27
type Promise struct {
// contains filtered or unexported fields
}
Promise 是一次性异步结果的生产者端。
Promise 可安全地由多个 goroutine 竞争完成;只有第一次 Resolve 成功。
type Result ¶ added in v0.4.27
Result 保存 Future 的一次产出值或错误。
type Scope ¶ added in v0.4.27
type Scope struct {
// contains filtered or unexported fields
}
Scope 将一组后台任务绑定到同一生命周期。
Scope 组合 Context 取消、关闭后拒绝注册、活动任务计数和完成信号。Close 不会 强制终止 goroutine;任务必须观察传入的 Context 才能及时退出。
func (*Scope) AsyncScope ¶ added in v0.4.27
AsyncScope 返回自身,使 Scope 可直接作为生命周期 Scope 提供者使用。
func (*Scope) Stats ¶ added in v0.4.27
func (scope *Scope) Stats() ScopeStats
Stats 返回 Scope 的并发安全统计快照。
type ScopeStats ¶ added in v0.4.27
type ScopeStats struct {
Spawned int64 // 成功注册过的任务总数。
Active int64 // 当前仍未退出的任务数。
Completed int64 // 正常返回的任务总数。
Canceled int64 // 在 Context 已取消状态下退出的任务总数。
Rejected int64 // Scope 关闭后拒绝的任务总数。
Closed bool // Scope 是否已经关闭。
}
ScopeStats 是异步作用域的瞬时统计快照。
type Signal ¶ added in v0.4.27
type Signal struct {
// contains filtered or unexported fields
}
Signal 是可重放、无结果的完成通知。
func CompletedSignal ¶ added in v0.4.27
func CompletedSignal() Signal
CompletedSignal 创建已经完成的 Signal。
type Stream ¶ added in v0.4.27
type Stream struct {
// contains filtered or unexported fields
}
Stream 是连续 Result 的只读消费者视图。
Stream 采用单消费语义;多个消费者读取 Chan 会竞争元素。需要广播时应在更高层 使用事件或消息总线,而不是共享 Stream。
type Subscription ¶ added in v0.4.27
type Subscription struct {
// contains filtered or unexported fields
}
Subscription 表示一个可取消的 Future 完成订阅。
func (Subscription) Cancel ¶ added in v0.4.27
func (subscription Subscription) Cancel() bool
Cancel 取消尚未执行的订阅。成功移除返回 true;已完成、已取消或零值返回 false。