一、前置思考
多设备协同场景中,事件传递是最基础的需求——手机端点击"分享"按钮,平板端要感知到;PC端修改了文档标题,智慧屏端要更新标题展示。简单地用KVStore轮询监听效率极低,分布式事件总线(Distributed Event Bus)提供了高效的发布-订阅模式。
本文聚焦:
- 分布式事件总线的架构设计
- 事件全局唯一性与顺序保证
- 延迟解耦(Deferred Event)与重放机制
- 事件的TTL生命周期管理
二、核心原理
2.1 事件总线架构
发布者 (Publisher) 订阅者 (Subscriber) │ │ ├─ publish(e:Event) ──┐ │ │ ▼ │ │ ┌──────────┐ │ │ │ EventBus │ │ │ │ ┌──────┐ │ │ │ │ │Router│──├───┤─→ Topic匹配 │ │ └──────┘ │ │ │ │ ┌──────┐ │ │ │ │ │Queue │ │ │ 先入先出 │ │ └──────┘ │ │ │ │ ┌──────┐ │ │ │ │ │Store │ │ │ 持久化 │ │ └──────┘ │ │ │ └──────────┘ │ │ │ ▼ ▼ 通过软总线跨设备分发 onEvent(Topic)回调2.2 事件定义
interfaceDistributedEvent{eventId:string;// 全局唯一ID (UUID v4)topic:string;// 事件主题 (如: "doc:update")sourceDeviceId:string;// 源设备IDsourceAppId:string;// 源应用IDtimestamp:number;// 事件发生时间 (UTC毫秒)ttl:number;// 生存时间(ms),超时丢弃priority:number;// 优先级 0-10payload:string;// 事件载荷(JSON)sequenceNumber:number;// 全局递增序号}// 事件生成器classEventFactory{privatesequenceCounter:number=0;privatedeviceId:string;constructor(deviceId:string){this.deviceId=deviceId;}createEvent(topic:string,payload:string,priority:number=5,ttl:number=30000):DistributedEvent{this.sequenceCounter++;return{eventId:this.generateUUID(),topic:topic,sourceDeviceId:this.deviceId,sourceAppId:'com.example.app',timestamp:Date.now(),ttl:ttl,priority:priority,payload:payload,sequenceNumber:this.sequenceCounter};}privategenerateUUID():string{// 简化UUID生成returnthis.deviceId+'-'+Date.now()+'-'+Math.random().toString(36).slice(2,10);}}2.3 分布式事件总线核心实现
classDistributedEventBus{privatekvStore:distributedKVStore.SingleKVStore|null=null;privatesubscribers:Map<string,SubscriberInfo[]>=newMap();privateeventQueue:DistributedEvent[]=[];privatereadonlyMAX_QUEUE_SIZE:number=1000;privatereadonlyEVENT_KEY_PREFIX:string='evt:';// 发布事件asyncpublish(event:DistributedEvent):Promise<void>{if(this.kvStore===null)return;// 入队(本地队列+KVStore)this.enqueue(event);// 写入KVStore触发远端同步constkey:string=this.EVENT_KEY_PREFIX+event.eventId;constvalue:string=JSON.stringify(event);awaitthis.kvStore.put(key,value);awaitthis.kvStore.sync([],distributedKVStore.SyncMode.PUSH_ONLY);// 设置TTL自动清理setTimeout(()=>{this.cleanupEvent(event.eventId);},event.ttl);}// 订阅subscribe(topic:string,callback:(event:DistributedEvent)=>void):string{constsubscriberId:string=topic+'-'+Date.now();letsubs:SubscriberInfo[]|undefined=this.subscribers.get(topic);if(subs===undefined){subs=[];this.subscribers.set(topic,subs);}subs.push({id:subscriberId,callback:callback,topic:topic});returnsubscriberId;}// 注销unsubscribe(subscriberId:string):void{constentries:MapIterator<[string,SubscriberInfo[]]>=this.subscribers.entries();for(letentry=entries.next();!entry.done;entry=entries.next()){consttopic:string=entry.value[0];constsubs:SubscriberInfo[]=entry.value[1];constnewSubs:SubscriberInfo[]=[];for(leti:number=0;i<subs.length;i++){if(subs[i].id!==subscriberId){newSubs.push(subs[i]);}}this.subscribers.set(topic,newSubs);}}// 处理远端事件privateonRemoteEvent(event:DistributedEvent):void{// 检查是否过期if(Date.now()-event.timestamp>event.ttl)return;// 匹配订阅者constsubs:SubscriberInfo[]|undefined=this.subscribers.get(event.topic);if(subs!==undefined){for(leti:number=0;i<subs.length;i++){subs[i].callback(event);}}}privateenqueue(event:DistributedEvent):void{this.eventQueue.push(event);// 按优先队列序排列this.eventQueue.sort((a:DistributedEvent,b:DistributedEvent)=>{if(a.priority!==b.priority)returnb.priority-a.priority;returna.sequenceNumber-b.sequenceNumber;});// 队列容量限制if(this.eventQueue.length>this.MAX_QUEUE_SIZE){this.eventQueue.shift();}}privateasynccleanupEvent(eventId:string):Promise<void>{if(this.kvStore!==null){awaitthis.kvStore.delete(this.EVENT_KEY_PREFIX+eventId);}}}interfaceSubscriberInfo{id:string;topic:string;callback:(event:DistributedEvent)=>void;}三、延迟解耦模式
// 离线设备事件延迟投递classDeferredEventDelivery{privatependingEvents:Map<string,DistributedEvent[]>=newMap();// 事件发布时目标设备离线 → 加入pendingdeferEvent(deviceId:string,event:DistributedEvent):void{letqueue:DistributedEvent[]|undefined=this.pendingEvents.get(deviceId);if(queue===undefined){queue=[];this.pendingEvents.set(deviceId,queue);}queue.push(event);// 限制pending队列大小if(queue.length>100)queue.shift();}// 设备上线 → 批量投递pending事件deliverDeferredEvents(deviceId:string):void{constqueue:DistributedEvent[]|undefined=this.pendingEvents.get(deviceId);if(queue===undefined||queue.length===0)return;console.info('[EventBus] 投递'+String(queue.length)+'个延迟事件到'+deviceId);// 按顺序投递for(leti:number=0;i<queue.length;i++){// 重新发布事件this.retryPublish(queue[i]);}this.pendingEvents.delete(deviceId);}}四、避坑速查
| 坑 | 现象 | 原因 | 解决 |
|---|---|---|---|
| 事件丢失 | 订阅者收不到事件 | KVStore的event key被过早清理 | TTL至少设为30s |
| 事件重复 | 订阅者收到重复事件 | 网络重传导致 | 订阅者用eventId去重 |
| 事件乱序 | 处理顺序与发送顺序不一致 | 网络延迟差异 | sequenceNumber排序本地重排 |
| 队列溢出 | 高频事件导致内存飙升 | 无队列上限 | 限制MAX_QUEUE_SIZE=1000,淘汰旧事件 |
| 订阅泄漏 | 关闭页面后仍在收事件 | 未unsubscribe | aboutToDisappear中注销订阅 |
五、总结
分布式事件总线设计要点:
- 全局唯一eventId + sequenceNumber保证顺序
- TTL自动过期,防止KVStore膨胀
- 延迟投递支持离线设备
- 优先级队列确保关键事件优先处理