Skip to content

Fast.EventBus

逐成员 API、参数与返回参考

版本 3.5.28;目标 net8.0net9.0net10.0;依赖 Fast.Runtime。默认实现是进程内、有界 Channel,不提供持久化或跨进程投递保证。

bash
dotnet add package Fast.EventBus

注册和订阅

services.AddEventBus() 返回服务集合。默认通道容量 3000,注册单例发布者、动态订阅工厂及扫描到的 IEventSubscriber,并注册后台消费服务。

csharp
using Fast.EventBus;
using Microsoft.AspNetCore.Builder;
using Microsoft.AspNetCore.Http;
using Microsoft.Extensions.DependencyInjection;

var builder = WebApplication.CreateBuilder(args);
builder.Services.AddEventBus();
var app = builder.Build();
app.MapPost("/notices", async (Notice notice, IEventPublisher publisher, CancellationToken cancellationToken) =>
{
    await publisher.PublishAsync("notice.received", notice, cancellationToken);
    return Results.Accepted();
});
app.Run();

public sealed record Notice(string Message);

public sealed class NoticeSubscriber : IEventSubscriber
{
    [EventSubscribe("notice.received")]
    public Task ReceiveAsync(EventHandlerExecutingContext context)
    {
        if (context.Source.Payload is Notice notice)
        {
            Console.WriteLine(notice.Message);
        }
        return Task.CompletedTask;
    }
}

订阅方法使用 Task Method(EventHandlerExecutingContext) 形式,类型须被框架扫描发现。订阅者是单例,避免捕获请求 Scoped 服务或将可变请求状态保存在实例中。

该示例是消费端 Program.cs,使用 Microsoft.NET.Sdk.Web 和本页支持的 .NET 8/9/10。HTTP 202 仅表示入队完成,不能当作业务处理完成。PublishDelayAsync("notice.received", 1000, notice, cancellationToken) 可延迟 1 秒入队;这不是持久化定时任务。文档维护不会启动此宿主。

发布与动态订阅 API

入口语义
IEventPublisher.PublishAsync(IEventSource)等待消息入队;不代表所有订阅处理完成
PublishAsync(string/Enum eventId, object payload = null, CancellationToken cancellationToken = default)构造消息并入队
PublishDelayAsync(IEventSource, long delay)延迟单位毫秒,等待延迟和入队;负值抛 ArgumentOutOfRangeException
PublishDelayAsync(string/Enum eventId, long delay, object payload = null, CancellationToken cancellationToken = default)同上;取消/写入异常可从返回任务观察
IEventBusFactory.Subscribe(...)接受事件 ID、Func<EventHandlerExecutingContext, Task>、可选特性/MethodInfo/取消令牌
IEventBusFactory.Unsubscribe(string eventId, CancellationToken cancellationToken = default)提交移除订阅操作;返回 Task

ChannelEventSource 包含 EventIdPayloadCreatedTimeCancellationToken;创建时间默认 DateTime.UtcNowEventBusToString/EventBusToEnum 用于框架的枚举事件标识转换,不能自行假设仅等于枚举数值字符串。

重试和处理上下文

EventSubscribeAttribute 支持 string/Enum 标识,其他类型抛 ArgumentException。默认 NumRetries = 0RetryTimeout = 1000 毫秒、Order = 0GCCollect = false;可指定 ExceptionTypesFallbackPolicy

EventHandlerContext 提供 SourcePropertiesHandlerMethodAttribute;执行前后上下文分别带时间,完成上下文还有执行异常。IEventHandlerMonitor 提供 OnExecutingAsync/OnExecutedAsyncIEventFallbackPolicy.CallbackAsync 接收执行上下文和异常。

重试会重复执行处理逻辑,业务需自行保证幂等。多个处理器可能并发执行,Order 不构成串行事务保证。进程退出会丢失尚未完成的内存消息;不要用于要求可靠持久投递的任务。

来源与验证

依据 发布接口延迟实现。本轮未编译片段、未启动后台消费服务。

嵌套发布与有界背压

普通发布者在默认有界队列满时等待空位。处理器正在消费同一队列时再次发布:有空位就入队;满队列立即抛出明确的 InvalidOperationException,而不是等待自己释放容量。不改成无限队列、不静默丢事件、不启动脱离生命周期的 Task.Run。调用方需在当前处理完成后发布,或显式处理容量不足。

只扫描闭合、非抽象、可实例化的订阅者、监视器和降级策略。显式注册优先;多个自动发现的单一监视器需要调用方明确选择。