Skip to content
bitzorcas
中EN

Tutorial

接入实时通知能力:Server-Sent Events (SSE) 高性能流式推送

遵循“持久化优先,推送第二”的架构铁律,使用 Server-Sent Events (SSE) 实现律所立案审批待办、利冲预警的高性能流式长连接推送与前端断线重连。

Last updated

在许多企业级系统的早期演进中,开发团队经常犯下一个致命错误:把推送通道当成了数据存储。为了追求“快”,后端直接通过内存消息队列或原生 WebSocket 将审批通知推向客户端,而跳过了数据库持久化。一旦用户切入后台、网络产生短暂抖动,或服务端实例发生热重载,未持久化的关键预警便彻底湮灭,导致主任律师错失利益冲突阻断或开庭时限提醒,引发重大合规执业事故。

BitzOrcas.Modern 坚定确立“持久化优先,推送第二(Persistence First, Streaming Second)”的架构铁律。所有业务通知必须首先作为实体持久化至数据库,生成唯一审计流水与收件箱状态;随后通过基于 HTTP/2 的 Server-Sent Events (SSE) 流式通道分发至在线用户。即使客户端处于离线状态,重新连线后亦可秒级拉取完整的持久化收件箱。

本教程将以律所立案待办与利益冲突实时警报为业务场景,带你从零实现一套高可用、抗重放且支持租户隔离的实时通知流。

实时通知流水线架构时序

"SQL Server 通知台账表""INotificationRepository""CAP 事务发件箱消费者""租户通知通道中心 (NotificationChannelHub)""SSE 流式端点 (/api/notifications/stream)""SQL Server 通知台账表""INotificationRepository""CAP 事务发件箱消费者""租户通知通道中心 (NotificationChannelHub)""SSE 流式端点 (/api/notifications/stream)"维持 HTTP/2 流式通道 (每 15 秒发送 :ping 心跳保活)步骤 1:持久化优先守门步骤 2:流式通道增量推送"主任合伙人前端 (Browser SPA)"1. 建立 SSE 长连接 (GET /api/notifications/stream, 带 JWT)12. 注册当前租户/用户的长连接通道 (Channel<NotificationEvent>)23. 消费到立案审批事件,构造 Notification 聚合根34. SaveAsync 写入数据库,持久化收件箱事实45. 触发推送 PublishToUserAsync(tenantId, userId, payload)56. 写入内存 Channel,触发 SSE StreamWriter67. 推送 SSE 数据包 (event: matter_alert\ndata: {...}\n\n)7"主任合伙人前端 (Browser SPA)"

架构权衡:为什么选择 SSE 而非 WebSocket?

在企业级管理端应用中,绝大多数实时交互均为服务端向客户端的单向状态广播(如待办提醒、报表生成进度、合规警报):

评估维度Server-Sent Events (SSE)WebSocketBitzOrcas 架构考量
底层协议标准 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>,构建内存级租户隔离的分发通道:

src/Platform/Notifications/BitzOrcas.Platform.Notifications.Infrastructure/Channels/NotificationChannelHub.cs
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 响应体:

src/Platform/Notifications/BitzOrcas.Platform.Notifications.Application/Endpoints/NotificationStreamEndpoint.cs
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 库)接入流式推送:

frontend/packages/platform-sdk/src/notifications/useNotificationStream.ts
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]);
}

核心治理小结

本方案严密实现了以下工程优势:

  1. 零消息丢失:持久化在先,即使推送失败或前端未在线,用户再次打开系统时仍可通过 /api/notifications/inbox 查询完整的历史信封;
  2. 极高并发密度:单台 API Host 借助异步 Channel<T> 与 ValueTask,可承载数万个并发保活连接,且内存开销极低;
  3. 天然穿透网关:完全兼容 YARP 反向代理与 HTTP/2 多路复用,零额外运维负担。

100%

滚轮或按钮缩放 · 放大后拖动画面 · 双击切换 100% / 200%