构建万级并发实时数据发布网需解耦通知契约并下沉至异步基础设施:通过eventconsumer、deliverystrategy、eventrouter接口多态实现可插拔通知;注册改用分布式中心存元数据,通知转为异步任务分发;引入响应式流控制背压;按主题分片与消费者分组双维度横向扩展。

要构建支持每秒万级并发的实时数据发布网,单靠观察者模式基础结构远远不够——它本质是同步、阻塞、内存内的一对多通知机制,原生无法承载高并发压力。真正可行的路径,是用接口多态“解耦通知契约”,再将观察者模式退为逻辑模型,把实际通知行为下沉到异步、分发、可伸缩的基础设施层。
用接口多态定义统一但可插拔的通知契约
不把 Observer 硬编码为具体实现类,而是定义细粒度、职责单一的接口:
-
EventConsumer:只声明
onEvent(DataEvent event),所有下游系统(前端 WebSocket、风控服务、日志聚合)都实现它,但各自处理逻辑隔离 -
DeliveryStrategy:声明
deliver(EventConsumer consumer, DataEvent event),允许按场景切换策略——比如对高优先级用户走直连推送,对批量报表走 Kafka 批量写入 - EventRouter:根据 event.type、tenantId、region 等字段路由到不同 DeliveryStrategy 实例,实现流量隔离与灰度能力
观察者注册不再存 List,而存轻量元数据
传统 ArrayList 存一万观察者,notify() 遍历本身就成了瓶颈。应改为:
- 注册时只存
consumerId + topic + deliveryPolicy到分布式注册中心(如 etcd 或 Redis Sorted Set),不加载实例 - 主题状态变更时,不调用 update(),而是生成一个
EventPublishTask,交由任务队列(如 RocketMQ 延迟队列或自研分片任务调度器)异步分发 - 每个工作节点拉取属于本分片的消费者列表,批量构造事件并投递,避免单点串行遍历
用响应式流替代同步回调,控制背压与吞吐
Java 生态中,直接让 Observer.update() 同步执行会阻塞主线程。正确做法是:
- 具体观察者实现
EventConsumer时,内部封装成Flux<dataevent></dataevent>或Flow.Publisher - DeliveryStrategy 调用
consumer.asPublisher().subscribe(subscriber),由 Project Reactor 或 JDK Flow 自动协调缓冲、丢弃、重试策略 - 当某下游消费慢时,上游自动降速(背压),而非堆积内存或抛异常——这对万级并发下的稳定性至关重要
横向扩展靠“主题分片 + 消费者分组”双维度拆分
万级并发不是靠单机扛,而是靠分布协同:
- 主题按业务域分片:user_profile_update → shard-01 ~ shard-08;每个分片部署独立 Publisher 实例
- 观察者按消费能力分组:实时看板组(要求 sub-50ms)、离线分析组(容忍分钟级延迟)、审计归档组(仅需 at-least-once)
- 同一事件可同时投递到多组,但每组使用独立线程池与限流器,互不影响











