如何让C#发布订阅模式中订阅者的通知处理函数改为异步
要把OnNotificationReceived改成异步方法,需要从委托定义、事件触发、方法实现三个核心部分调整,具体修改如下:
1. 修改发布者的委托与事件定义
异步事件处理不能使用void返回的委托(会导致异步操作无法被等待、异常无法捕获),需要改为返回Task的委托:
class Publisher{ // ...其他属性不变 // 把void改为Task,支持异步回调 public delegate Task Notify(Publisher p, NotificationEvent e); // 事件类型对应更新 public event Notify OnPublish; // ...构造函数不变 }
2. 改造发布者的Publish方法为异步
原来的Thread.Sleep会阻塞线程,改用await Task.Delay实现非阻塞等待,同时调整方法为异步,并处理多播委托的异步等待:
class Publisher{ // ...其他代码不变 // 改为async Task方法,支持取消终止 public async Task Publish(CancellationToken cancellationToken = default) { while (!cancellationToken.IsCancellationRequested) { // 非阻塞等待指定间隔 await Task.Delay(NotificationInterval, cancellationToken); // 线程安全获取事件委托副本(避免触发时订阅者变更) var publishEvent = OnPublish; if (publishEvent == null) continue; var notificationObj = new NotificationEvent(DateTime.Now, "New Notification Arrived from"); // 拆分多播委托,等待所有订阅者的异步操作完成 var handlers = publishEvent.GetInvocationList().Cast<Notify>(); await Task.WhenAll(handlers.Select(handler => handler(this, notificationObj))); } } }
3. 修改订阅者的异步处理方法
把OnNotificationReceived改为async Task的虚拟方法,添加async关键字并修改返回类型,内部可直接编写异步逻辑:
class Subscriber{ // ...其他属性、构造函数不变 // 异步虚拟方法 protected virtual async Task OnNotificationReceived(Publisher p, NotificationEvent e) { // 示例异步操作:比如等待IO、调用远程API等 await Task.Delay(100); Console.WriteLine($"Hey {SubscriberName}, {e.NotificationMessage} - {p.PublisherName} at {e.NotificationDate:yyyy-MM-dd HH:mm:ss}"); } }
4. 订阅/取消订阅逻辑无需修改
因为方法签名(返回Task、参数一致)匹配新的委托,原有的Subscribe和Unsubscribe代码可以直接保留:
class Subscriber{ // ...其他代码不变 public void Subscribe(Publisher p){ p.OnPublish += OnNotificationReceived; } public void Unsubscribe(Publisher p){ p.OnPublish -= OnNotificationReceived; } }
完整可运行代码
补充NotificationEvent类定义
public class NotificationEvent { public DateTime NotificationDate { get; } public string NotificationMessage { get; } public NotificationEvent(DateTime date, string message) { NotificationDate = date; NotificationMessage = message; } }
修改后的Publisher类
class Publisher{ public string PublisherName { get; private set; } public int NotificationInterval { get; private set; } public delegate Task Notify(Publisher p, NotificationEvent e); public event Notify OnPublish; public Publisher(string _publisherName, int _notificationInterval){ PublisherName = _publisherName; NotificationInterval = _notificationInterval; } public async Task Publish(CancellationToken cancellationToken = default){ while (!cancellationToken.IsCancellationRequested){ await Task.Delay(NotificationInterval, cancellationToken); var publishEvent = OnPublish; if (publishEvent == null) continue; NotificationEvent notificationObj = new NotificationEvent(DateTime.Now, "New Notification Arrived from"); var handlers = publishEvent.GetInvocationList().Cast<Notify>(); await Task.WhenAll(handlers.Select(h => h(this, notificationObj))); } } }
修改后的Subscriber类
class Subscriber{ public string SubscriberName { get; private set; } public Subscriber(string _subscriberName){ SubscriberName = _subscriberName; } public void Subscribe(Publisher p){ p.OnPublish += OnNotificationReceived; } public void Unsubscribe(Publisher p){ p.OnPublish -= OnNotificationReceived; } protected virtual async Task OnNotificationReceived(Publisher p, NotificationEvent e){ await Task.Delay(50); Console.WriteLine($"Hey {SubscriberName}, {e.NotificationMessage} - {p.PublisherName} at {e.NotificationDate:yyyy-MM-dd HH:mm:ss}"); } }
使用示例
var cts = new CancellationTokenSource(); var publisher = new Publisher("TechNews", 2000); var alice = new Subscriber("Alice"); var bob = new Subscriber("Bob"); alice.Subscribe(publisher); bob.Subscribe(publisher); // 启动发布任务 _ = publisher.Publish(cts.Token); Console.WriteLine("Press any key to stop..."); Console.ReadKey(); cts.Cancel();
内容的提问来源于stack exchange,提问作者Razvan Emil
相关产品推荐
相关产品推荐

