You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何让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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.18 14:45:34