在许多企业级系统的早期演进中,开发团队经常犯下一个致命错误:把推送通道当成了数据存储。为了追求“快”,后端直接通过内存消息队列或原生 WebSocket 将审批通知推向客户端,而跳过了数据库持久化。一旦用户切入后台、网络产生短暂抖动,或服务端实例发生热重载,未持久化的关键预警便彻底湮灭,导致主任律师错失利益冲突阻断或开庭时限提醒,引发重大合规执业事故。
BitzOrcas.Modern 坚定确立“持久化优先,推送第二(Persistence First, Streaming Second)”的架构铁律。所有业务通知必须首先作为实体持久化至数据库,生成唯一审计流水与收件箱状态;随后通过基于 HTTP/2 的 Server-Sent Events (SSE) 流式通道分发至在线用户。即使客户端处于离线状态,重新连线后亦可秒级拉取完整的持久化收件箱。
本教程将以律所立案待办与利益冲突实时警报为业务场景,带你从零实现一套高可用、抗重放且支持租户隔离的实时通知流。
实时通知流水线架构时序
架构权衡:为什么选择 SSE 而非 WebSocket?
在企业级管理端应用中,绝大多数实时交互均为服务端向客户端的单向状态广播(如待办提醒、报表生成进度、合规警报):
| 评估维度 | Server-Sent Events (SSE) | WebSocket | BitzOrcas 架构考量 |
|---|---|---|---|
| 底层协议 | 标准 HTTP/1.1 与 HTTP/2 | 独立的 ws:// / wss:// 协议升级 | SSE 与现有 ASP.NET Core 10 级安全管道 100% 融合,无协议切换成本 |
| 反向代理穿透 | 原生兼容 YARP、Nginx、Cloudflare | 需特殊配置 Connection: Upgrade | 零额外代理配置,通过 YARP 网关时天然继承分布式链路上下文 |
| 断线自动重连 | 浏览器原生 EventSource 自动重连 | 需前端手工编写心跳与指数退避代码 | 极大降低前端状态机复杂度,网络恢复后原生秒级自愈 |
| 防火墙友好度 | 走标准 443 端口,极少被安全策略拦截 | 经常被严格的企业内网安全代理封锁 | 适合银行、律所等高安全级别企业内网部署 |
第一步:构建租户隔离的消息通道中心(Channel Hub)
利用 .NET 高性能无锁并发原语 System.Threading.Channels.Channel<T>,构建内存级租户隔离的分发通道:
using System.Collections.Concurrent;using System.Threading.Channels;
namespace BitzOrcas.Platform.Notifications.Infrastructure.Channels;
/// <summary>/// 租户级实时通知流式分发通道中心/// </summary>public sealed class NotificationChannelHub{ // 双层字典保证租户数据隔离:TenantId -> (UserId -> Channel) private readonly ConcurrentDictionary<string, ConcurrentDictionary<string, Channel<NotificationPayload>>> _tenantChannels = new();
/// <summary> /// 为特定用户注册并获取流式读取通道 /// </summary> public ChannelReader<NotificationPayload> Subscribe(string tenantId, string userId) { var userMap = _tenantChannels.GetOrAdd(tenantId, _ => new ConcurrentDictionary<string, Channel<NotificationPayload>>());
// 使用单读单写的高吞吐 bounded 内存通道,设置背压缓冲容量为 100 var channel = userMap.GetOrAdd(userId, _ => Channel.CreateBounded<NotificationPayload>(new BoundedChannelOptions(100) { FullMode = BoundedChannelFullMode.DropOldest }));
return channel.Reader; }
/// <summary> /// 注销用户长连接会话 /// </summary> public void Unsubscribe(string tenantId, string userId) { if (_tenantChannels.TryGetValue(tenantId, out var userMap)) { userMap.TryRemove(userId, out _); } }
/// <summary> /// 向指定租户下的目标用户推送增量通知 /// </summary> public async ValueTask PublishToUserAsync(string tenantId, string userId, NotificationPayload payload, CancellationToken cancellationToken = default) { if (_tenantChannels.TryGetValue(tenantId, out var userMap) && userMap.TryGetValue(userId, out var channel)) { // 写入通道,若用户在线立即触发 HTTP 流式刷新 await channel.Writer.WriteAsync(payload, cancellationToken); } }}
public sealed record NotificationPayload( string NotificationId, string EventType, string Title, string Content, DateTimeOffset CreatedAtUtc);第二步:编写高性能 SSE Minimal API 流端点
在通知模块中公开流端点,直接写入 HTTP 响应体:
using System.Text.Json;using BitzOrcas.Application.Abstractions.Tenancy;using BitzOrcas.Application.Security;using BitzOrcas.Platform.Notifications.Infrastructure.Channels;using Microsoft.AspNetCore.Builder;using Microsoft.AspNetCore.Http;using Microsoft.AspNetCore.Routing;
namespace BitzOrcas.Platform.Notifications.Application.Endpoints;
public static class NotificationStreamEndpoint{ public static void MapNotificationStream(this IEndpointRouteBuilder app) { app.MapGet("/api/notifications/stream", async ( HttpContext context, ICurrentUser currentUser, ICurrentTenant currentTenant, NotificationChannelHub channelHub, CancellationToken cancellationToken) => { // 1. 设置 Server-Sent Events 标准协议响应标头 context.Response.Headers.ContentType = "text/event-stream"; context.Response.Headers.CacheControl = "no-cache"; context.Response.Headers.Connection = "keep-alive";
var tenantId = currentTenant.Tenant.EffectiveTenantId; var userId = currentUser.UserId; var reader = channelHub.Subscribe(tenantId, userId);
try { // 发送连接就绪事件 await context.Response.WriteAsync("event: connected\ndata: {\"status\":\"ok\"}\n\n", cancellationToken); await context.Response.Body.FlushAsync(cancellationToken);
// 2. 循环消费通道中的实时消息流 while (!cancellationToken.IsCancellationRequested) { // 设置 15 秒心跳保活检测,防止云厂商 NAT 网关或反向代理单向超时断开 using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); cts.CancelAfter(TimeSpan.FromSeconds(15));
try { var payload = await reader.ReadAsync(cts.Token); var json = JsonSerializer.Serialize(payload);
// 格式化为 SSE 标准包体:event: <name>\ndata: <json>\n\n await context.Response.WriteAsync($"event: {payload.EventType}\n", cancellationToken); await context.Response.WriteAsync($"data: {json}\n\n", cancellationToken); await context.Response.Body.FlushAsync(cancellationToken); } catch (OperationCanceledException) when (!cancellationToken.IsCancellationRequested) { // 15 秒内无业务消息,推送 SSE 空白注释包保活心跳 await context.Response.WriteAsync(":ping\n\n", cancellationToken); await context.Response.Body.FlushAsync(cancellationToken); } } } finally { // 3. 连接断开时确保资源回收注销 channelHub.Unsubscribe(tenantId, userId); } }) .RequireAuthorization() .WithTags("Notifications") .WithSummary("建立租户级实时通知 SSE 流式长连接"); }}第三步:前端 TypeScript 原生客户端集成
在前端应用中,使用浏览器原生的 EventSource(或支持携带自定义 Header 的 fetch-event-source 库)接入流式推送:
import { useEffect } from 'react';
export function useNotificationStream(accessToken: string) { useEffect(() => { if (!accessToken) return;
// 建立基于标准 HTTP 的 SSE 长连接 (通过 URL 参数或 Cookie 传递鉴权凭证) const eventSource = new EventSource(`/api/notifications/stream?access_token=${encodeURIComponent(accessToken)}`);
// 监听立案审批提醒事件 eventSource.addEventListener('matter_alert', (event: MessageEvent) => { const data = JSON.parse(event.data); console.info('[SSE] 收到待办与利益冲突预警:', data);
// 触发 UI 侧全局 Toast 提醒并增量刷新未读数角标 window.dispatchEvent(new CustomEvent('bitz:notification:received', { detail: data })); });
// 原生处理网络故障与自动重连 eventSource.onerror = (err) => { console.warn('[SSE] 连接中断,浏览器将自动根据指数退避策略重试...', err); };
return () => { eventSource.close(); }; }, [accessToken]);}核心治理小结
本方案严密实现了以下工程优势:
- 零消息丢失:持久化在先,即使推送失败或前端未在线,用户再次打开系统时仍可通过
/api/notifications/inbox查询完整的历史信封; - 极高并发密度:单台 API Host 借助异步
Channel<T>与ValueTask,可承载数万个并发保活连接,且内存开销极低; - 天然穿透网关:完全兼容 YARP 反向代理与 HTTP/2 多路复用,零额外运维负担。