Documentation
¶
Overview ¶
Package event 提供基于 Actor 的进程内事件派发能力。
Package event 提供基于 Actor 的进程内事件派发能力。
Package event 提供基于 Actor 的进程内事件派发能力。
Package event 提供基于 Actor 的进程内事件派发能力。
Package event 提供基于 Actor 的进程内事件派发能力。
Index ¶
- Variables
- type Codec
- type Config
- type DeadLetter
- type DeadLetterSink
- type DeadLetterSinkFunc
- type DeliveryStatus
- type Dispatcher
- type ErrorClassifier
- type Event
- type EventBus
- func (d *EventBus) Health() HealthSnapshot
- func (d *EventBus) Metrics() MetricsSnapshot
- func (d *EventBus) Publish(ctx context.Context, e Event) (PublishResult, error)
- func (d *EventBus) Shutdown(ctx context.Context) error
- func (d *EventBus) State() State
- func (d *EventBus) Subscribe(topic string, handler Handler, opts ...SubscribeOption) (Subscription, error)
- type Handler
- type HandlerFunc
- type HealthSnapshot
- type HealthStatus
- type MetricsSnapshot
- type Option
- type PublishResult
- type Publisher
- type QueuePolicy
- type State
- type Store
- type StoredEvent
- type SubscribeConfig
- type SubscribeOption
- type Subscription
Constants ¶
This section is empty.
Variables ¶
var ( // ErrInvalidConfig 表示配置为空、越界或相互冲突。 ErrInvalidConfig = errors.New("event: invalid config") // ErrInvalidEvent 表示事件缺少必需的 ID 或 Topic。 ErrInvalidEvent = errors.New("event: invalid event") // ErrNoSubscriber 表示目标 Topic 当前没有订阅者。 ErrNoSubscriber = errors.New("event: no subscriber") // ErrQueueFull 表示订阅队列已达到容量上限。 ErrQueueFull = errors.New("event: queue full") // ErrClosed 表示派发器或订阅已经关闭。 ErrClosed = errors.New("event: dispatcher closed") // ErrNotReady 表示运行时资源尚未准备好。 ErrNotReady = errors.New("event: dispatcher not ready") ErrActorUnavailable = errors.New("event: actor unavailable") // ErrCodec 表示编码或解码过程失败。 ErrCodec = errors.New("event: codec failure") )
Functions ¶
This section is empty.
Types ¶
type Codec ¶
type Codec interface {
// Encode 在事件进入订阅队列前编码 Payload。
Encode(context.Context, any) ([]byte, error)
// Decode 为一次 Handler 调用生成独立的 Payload。
Decode(context.Context, []byte) (any, error)
}
Codec 负责序列化和反序列化 Event.Payload。
派发器在入队前只编码一次,并为每次 Handler 调用独立解码。 Encode 返回后,调用方必须将返回的字节切片视为只读数据。
type Config ¶
type Config struct {
// QueueCapacity 是每个订阅分区的队列容量上限。
QueueCapacity int
// Partitions 是每个订阅创建的固定分区数。
Partitions int
// PublishPolicy 定义队列满时的发布行为。
PublishPolicy QueuePolicy
// MaxAttempts 是单个事件允许的最大处理次数。
MaxAttempts int
// HandlerTimeout 限制单次 Handler 调用时长,零表示不限制。
HandlerTimeout time.Duration
// RetryInitialBackoff 是第一次重试的退避时长。
RetryInitialBackoff time.Duration
// RetryMaxBackoff 是重试退避的最大时长。
RetryMaxBackoff time.Duration
// ErrorClassifier 判断 Handler 错误是否可重试。
ErrorClassifier ErrorClassifier
// DeadLetterSink 接收最终失败或不可重试的事件。
DeadLetterSink DeadLetterSink
// ActorSystemName 是底层 Actor 系统的名称前缀。
ActorSystemName string
// Codec 定义发布编码和 Handler 解码策略。
Codec Codec
// Store 预留外部持久化适配边界,首期不会主动调用。
Store Store
}
Config 是 EventBus 的默认运行配置。
type DeadLetter ¶
type DeadLetter struct {
// Event 是失败事件的副本。
Event Event
// LastError 是最后一次处理错误。
LastError error
// Attempts 是已经执行的处理次数。
Attempts int
// Panicked 表示失败是否由 Handler panic 引起。
Panicked bool
// PanicValue 保存捕获到的 panic 值。
PanicValue any
// FailedAt 是进入死信时的时间。
FailedAt time.Time
}
DeadLetter 保存无法成功处理的事件及其最终失败信息。
type DeadLetterSink ¶
type DeadLetterSink interface {
// Handle 处理一条死信记录。
Handle(context.Context, DeadLetter) error
}
DeadLetterSink 在重试耗尽或错误不可重试时接收死信。
type DeadLetterSinkFunc ¶
type DeadLetterSinkFunc func(context.Context, DeadLetter) error
DeadLetterSinkFunc 将普通函数适配为 DeadLetterSink。
func (DeadLetterSinkFunc) Handle ¶
func (f DeadLetterSinkFunc) Handle(ctx context.Context, dead DeadLetter) error
Handle 调用底层函数处理死信。
type DeliveryStatus ¶
type DeliveryStatus uint8
DeliveryStatus 表示一次发布的总体投递状态。
const ( // DeliveryAccepted 表示事件已被所有目标订阅接收。 DeliveryAccepted DeliveryStatus = iota + 1 // DeliveryRejected 表示至少一个目标订阅拒绝接收。 DeliveryRejected // DeliveryCompleted 表示事件已完成处理。 DeliveryCompleted // DeliveryFailed 表示事件处理失败但尚未进入最终状态。 DeliveryFailed // DeliveryDeadLettered 表示事件已进入死信流程。 DeliveryDeadLettered // DeliveryCancelled 表示发布因关闭或取消而终止。 DeliveryCancelled )
type Dispatcher ¶
type Dispatcher interface {
Publisher
Subscribe(string, Handler, ...SubscribeOption) (Subscription, error)
Shutdown(context.Context) error
State() State
Metrics() MetricsSnapshot
Health() HealthSnapshot
}
Dispatcher 是事件发布、订阅、生命周期和观测能力的完整接口。
type Event ¶
type Event struct {
// ID 是业务侧用于幂等和追踪的事件标识。
ID string
// Topic 是订阅匹配使用的主题。
Topic string
// Key 决定分区;相同 Key 在同一订阅内保持 FIFO。
Key string
// Payload 是业务负载;配置 Codec 后会按 Handler 独立解码。
Payload any
// Metadata 保存可选的字符串元数据。
Metadata map[string]string
// CreatedAt 是事件创建时间,用于延迟指标计算。
CreatedAt time.Time
}
Event 是传递给 Handler 的不可变事件信封。 未配置 Codec 时,Payload 由调用方保证在 Publish 返回后只读。
type EventBus ¶
type EventBus struct {
// contains filtered or unexported fields
}
EventBus 是事件派发器的默认实现。
func (*EventBus) Shutdown ¶
Shutdown 停止接收新工作,等待已接收工作排空,并关闭底层运行时。 若 Context 超时,则取消内部处理并返回 Context 错误。
func (*EventBus) Subscribe ¶
func (d *EventBus) Subscribe(topic string, handler Handler, opts ...SubscribeOption) (Subscription, error)
Subscribe 为 Topic 注册 Handler,并启动该订阅的所有分区。
type Handler ¶
type Handler interface {
// Handle 执行业务处理逻辑;返回错误时由派发器决定重试或进入死信。
Handle(context.Context, Event) error
}
Handler 消费一个事件并返回处理结果。
type HandlerFunc ¶
HandlerFunc 将普通函数适配为 Handler。
type HealthSnapshot ¶
type HealthSnapshot struct {
// Status 是当前健康状态。
Status HealthStatus
// Reason 是状态的简要原因。
Reason string
}
HealthSnapshot 是 Health 返回的健康状态快照。
type HealthStatus ¶
type HealthStatus uint8
HealthStatus 表示派发器对外提供服务的健康状态。
const ( // HealthReady 表示派发器可以接收新工作。 HealthReady HealthStatus = iota + 1 // HealthDraining 表示派发器正在排空已接收工作。 HealthDraining // HealthNotReady 表示派发器不可接收新工作。 HealthNotReady )
type MetricsSnapshot ¶
type MetricsSnapshot struct {
// Published 是通过事件校验并进入发布流程的次数。
Published uint64
// Accepted 是成功进入订阅队列的投递次数。
Accepted uint64
// QueueFull 是因队列满而拒绝的投递次数。
QueueFull uint64
// Processed 是 Handler 成功完成的次数。
Processed uint64
// Failed 是 Handler 返回错误或 panic 的次数。
Failed uint64
// Retried 是已安排重试的次数。
Retried uint64
// DeadLettered 是进入死信流程的次数。
DeadLettered uint64
// Panicked 是捕获到 Handler panic 的次数。
Panicked uint64
// QueueDepth 是当前所有订阅队列中的事件数量。
QueueDepth int
// InFlight 是当前正在执行 Handler 的事件数量。
InFlight int
// MaxLatency 是观测到的最大处理延迟。
MaxLatency time.Duration
}
MetricsSnapshot 是 Metrics 返回的原子指标快照。
type PublishResult ¶
type PublishResult struct {
// Accepted 是成功进入订阅队列的数量。
Accepted int
// Rejected 是未能进入订阅队列的数量。
Rejected int
// Status 是本次发布的总体状态。
Status DeliveryStatus
}
PublishResult 描述一次 Publish 对各订阅的入队结果。
type Publisher ¶
type Publisher interface {
// Publish 将事件投递到匹配的订阅队列。
Publish(context.Context, Event) (PublishResult, error)
}
Publisher 定义发布事件的最小接口。
type QueuePolicy ¶
type QueuePolicy uint8
QueuePolicy 定义队列满时 Publish 的处理策略。
const ( // QueueReject 立即返回 ErrQueueFull。 QueueReject QueuePolicy = iota // QueueBlock 等待队列空间,并受 Context 取消或截止时间约束。 QueueBlock )
type Store ¶
type Store interface {
// Save 持久化一个事件记录。
Save(context.Context, StoredEvent) error
// Load 根据事件 ID 加载一个事件记录。
Load(context.Context, string) (StoredEvent, error)
}
Store 描述可选的外部持久化集成边界。 首期派发器不会主动调用该接口。
type StoredEvent ¶
type StoredEvent struct {
// Event 保存事件信封。
Event Event
// EncodedPayload 保存 Codec 编码后的负载。
EncodedPayload []byte
}
StoredEvent 是编码事件未来接入持久化存储时使用的边界结构。
首期实现不会调用 Store,也不提供崩溃恢复;未来接入外部可靠适配器时 可以复用该结构。
type SubscribeConfig ¶
type SubscribeConfig struct {
// QueueCapacity 是该订阅每个分区的队列容量。
QueueCapacity int
// Partitions 是该订阅的固定分区数。
Partitions int
// PublishPolicy 定义该订阅队列满时的行为。
PublishPolicy QueuePolicy
// HandlerTimeout 限制该订阅 Handler 的单次处理时长。
HandlerTimeout time.Duration
// MaxAttempts 是该订阅的最大处理次数。
MaxAttempts int
// RetryInitialBackoff 是该订阅第一次重试的退避时长。
RetryInitialBackoff time.Duration
// RetryMaxBackoff 是该订阅重试退避的上限。
RetryMaxBackoff time.Duration
// ErrorClassifier 判断该订阅错误是否可重试。
ErrorClassifier ErrorClassifier
// DeadLetterSink 接收该订阅的最终失败事件。
DeadLetterSink DeadLetterSink
// Codec 定义该订阅的 Payload 解码方式。
Codec Codec
}
SubscribeConfig 是单个订阅可覆盖的运行配置。
type Subscription ¶
type Subscription interface {
// ID 返回订阅唯一标识。
ID() string
// Topic 返回订阅主题。
Topic() string
// Close 停止订阅并排空或取消其队列中的工作。
Close() error
}
Subscription 表示一个 Topic 订阅及其生命周期控制句柄。
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
examples
|
|
|
basic
command
basic 示例展示最小可运行的事件发布、订阅和优雅关闭流程。
|
basic 示例展示最小可运行的事件发布、订阅和优雅关闭流程。 |
|
internal
|
|
|
actorbridge
Package actorbridge 隔离外部 Actor 运行时与事件派发器实现。
|
Package actorbridge 隔离外部 Actor 运行时与事件派发器实现。 |
|
handler
Package handler 提供 Handler 调用、panic 隔离和错误分类工具。
|
Package handler 提供 Handler 调用、panic 隔离和错误分类工具。 |
|
queue
Package queue 提供固定容量、并发安全的泛型队列。
|
Package queue 提供固定容量、并发安全的泛型队列。 |
|
retry
Package retry 提供重试退避和延迟任务调度能力。
|
Package retry 提供重试退避和延迟任务调度能力。 |
|
wakeup
Package wakeup 提供合并 Actor 唤醒请求的轻量状态门。
|
Package wakeup 提供合并 Actor 唤醒请求的轻量状态门。 |