基于事件总线EventBus实现邮件推送功能

2024-08-27 20:28

本文主要是介绍基于事件总线EventBus实现邮件推送功能,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

什么是事件总线

事件总线是对发布-订阅模式的一种实现。它是一种集中式事件处理机制,允许不同的组件之间进行彼此通信而又不需要相互依赖,达到一种解耦的目的。 关于这个概念,网上有很多讲解的,这里我推荐一个讲的比较好的(事件总线知多少)

什么是RabbitMQ

RabbitMQ这个就不用说了,想必到家都知道。

粗糙流程图

简单来解释就是:

1、定义一个事件抽象类

public abstract class EventData{/// <summary>/// 唯一标识/// </summary>public string Unique { get; set; }/// <summary>/// 是否成功/// </summary>public bool Success { get; set; }/// <summary>/// 结果/// </summary>public string Result { get; set; }}

2、定义一个事件处理抽象类,以及对应的一个队列消息执行的一个记录、

public abstract class EventHandler<T> where T : EventData{public async Task Handler(T eventData){await BeginHandler(eventData.Unique);eventData = await ProcessingHandler(eventData);if (eventData.Success)await FinishHandler(eventData);}/// <summary>///  开始处理/// </summary>/// <param name="unique"></param>/// <returns></returns>protected abstract Task BeginHandler(string unique);/// <summary>/// 处理中/// </summary>/// <param name="eventData"></param>/// <returns></returns>protected abstract Task<T> ProcessingHandler(T eventData);/// <summary>/// 处理完成/// </summary>/// <param name="eventData"></param>/// <returns></returns>protected abstract Task FinishHandler(T eventData);}[Table("Sys_TaskRecord")]public class TaskRecord : Entity<long>{/// <summary>/// 任务类型/// </summary>public TaskRecordType TaskType { get; set; }/// <summary>/// 任务状态/// </summary>public int TaskStatu { get; set; }/// <summary>/// 任务值/// </summary>public string TaskValue { get; set; }/// <summary>/// 任务结果/// </summary>public string TaskResult { get; set; }/// <summary>/// 任务开始时间/// </summary>public DateTime TaskStartTime { get; set; }/// <summary>/// 任务完成时间/// </summary>public DateTime? TaskFinishTime { get; set; }/// <summary>/// 任务最后更新时间/// </summary>public DateTime? LastUpdateTime { get; set; }/// <summary>/// 任务名称/// </summary>public string TaskName { get; set; }/// <summary>/// 附加数据/// </summary>public string AdditionalData { get; set; }}

3、定义一个邮件事件消息类,继承自EventData,以及一个邮件处理的Hanler继承自EventHandler

public class EmailEventData:EventData{/// <summary>/// 邮件内容/// </summary>public string Body { get; set; }/// <summary>/// 接收者/// </summary>public string Reciver { get; set; }}public class CreateEmailHandler<T> : Core.EventBus.EventHandler<T> where T : EventData{private IEmailService emailService;private IUnitOfWork unitOfWork;private ITaskRecordService taskRecordService;public CreateEmailHandler(IEmailService emailService, IUnitOfWork unitOfWork, ITaskRecordService taskRecordService){this.emailService = emailService;this.unitOfWork = unitOfWork;this.taskRecordService = taskRecordService;}protected override async Task BeginHandler(string unique){await taskRecordService.UpdateRecordStatu(Convert.ToInt64(unique), (int)MqMessageStatu.Processing);await unitOfWork.CommitAsync();}protected override async Task<T> ProcessingHandler(T eventData){try{EmailEventData emailEventData = eventData as EmailEventData;await emailService.SendEmail(emailEventData.Reciver, emailEventData.Reciver, emailEventData.Body, "[闲蛋]收到一条留言");eventData.Success = true;}catch (Exception ex){await taskRecordService.UpdateRecordFailStatu(Convert.ToInt64(eventData.Unique), (int)MqMessageStatu.Fail,ex.Message);await unitOfWork.CommitAsync();eventData.Success = false;}return eventData;}protected override async Task FinishHandler(T eventData){await taskRecordService.UpdateRecordSuccessStatu(Convert.ToInt64(eventData.Unique), (int)MqMessageStatu.Finish,"");await unitOfWork.CommitAsync();}

 4、接着就是如何把事件消息和事件Hanler关联起来,那么我这里思路就是把EmailEventData的类型和CreateEmailHandler的类型先注册到字典里面,这样我就可以根据EmailEventData找到对应的处理程序了,找类型还不够,如何创建实例呢,这里就还需要把CreateEmailHandler注册到DI容器里面,这样就可以根据容器获取对象了,如下

  public void AddSub<T, TH>()where T : EventDatawhere TH : EventHandler<T>{Type eventDataType = typeof(T);Type handlerType = typeof(TH);if (!eventhandlers.ContainsKey(typeof(T)))eventhandlers.TryAdd(eventDataType, handlerType);_serviceDescriptors.AddScoped(handlerType);}
-------------------------------------------------------------------------------------------------------------------public Type FindEventType(string eventName){if (!eventTypes.ContainsKey(eventName))throw new ArgumentException(string.Format("eventTypes不存在类名{0}的key", eventName));return eventTypes[eventName];}
------------------------------------------------------------------------------------------------------------------------------------------------------------public object FindHandlerType(Type eventDataType){if (!eventhandlers.ContainsKey(eventDataType))throw new ArgumentException(string.Format("eventhandlers不存在类型{0}的key", eventDataType.FullName));var obj = _buildServiceProvider(_serviceDescriptors).GetService(eventhandlers[eventDataType]);return obj;}
----------------------------------------------------------------------------------------------------------------------------------private static IServiceCollection AddEventBusService(this IServiceCollection services){string exchangeName = ConfigureProvider.configuration.GetSection("EventBusOption:ExchangeName").Value;services.AddEventBus(Assembly.Load("XianDan.Application").GetTypes()).AddSubscribe<EmailEventData, CreateEmailHandler<EmailEventData>>(exchangeName, ExchangeType.Direct, BizKey.EmailQueueName);return services;}

5、发送消息,这里代码简单,就是简单的发送消息,这里用eventData.GetType().Name作为消息的RoutingKey,这样消费这就可以根据这个key调用FindEventType,然后找到对应的处理程序了

using (IModel channel = connection.CreateModel())
{string routeKey = eventData.GetType().Name;string message = JsonConvert.SerializeObject(eventData);byte[] body = Encoding.UTF8.GetBytes(message);channel.ExchangeDeclare(exchangeName, exchangeType, true, false, null);channel.QueueDeclare(queueName, true, false, false, null);channel.BasicPublish(exchangeName, routeKey, null, body);
}

6、订阅消息,核心的是这一段

  Type eventType = _eventBusManager.FindEventType(eventName);  var eventData = (T)JsonConvert.DeserializeObject(body, eventType);  EventHandler<T> eventHandler = _eventBusManager.FindHandlerType(eventType)  as       EventHandler<T>;

public void Subscribe<T, TH>(string exchangeName, string exchangeType, string queueName)where T : EventDatawhere TH : EventHandler<T>{try{_eventBusManager.AddSub<T, TH>();IModel channel = connection.CreateModel();channel.QueueDeclare(queueName, true, false, false, null);channel.ExchangeDeclare(exchangeName, exchangeType, true, false, null);channel.QueueBind(queueName, exchangeName, typeof(T).Name, null);var consumer = new EventingBasicConsumer(channel);consumer.Received += async (model, ea) =>{string eventName = ea.RoutingKey;byte[] resp = ea.Body.ToArray();string body = Encoding.UTF8.GetString(resp);try{Type eventType = _eventBusManager.FindEventType(eventName);var eventData = (T)JsonConvert.DeserializeObject(body, eventType);EventHandler<T> eventHandler = _eventBusManager.FindHandlerType(eventType) as EventHandler<T>;await eventHandler.Handler(eventData);}catch (Exception ex){LogUtils.LogError(ex, "EventBusRabbitMQ", ex.Message);}finally{channel.BasicAck(ea.DeliveryTag, false);}};channel.BasicConsume(queueName, autoAck: false, consumer: consumer);}catch (Exception ex){LogUtils.LogError(ex, "EventBusRabbitMQ.Subscribe", ex.Message);}}

注意,这里我使用的时候有个小坑,就是最开始是用using包裹这个IModel channel = connection.CreateModel();导致最后程序启动后无法收到消息,然后去rabbitmq的管理界面发现没有channel连接,队列也没有消费者,最后发现可能是using执行完后就释放掉了,把using去掉就好了。

好了,到此,我的思路大概讲完了,现在我的网站留言也可以收到邮件了,那么多测试邮件,哈哈哈哈哈

文章转载自:灬丶

原文链接:https://www.cnblogs.com/MrHanBlog/p/18381572

体验地址:引迈 - JNPF快速开发平台_低代码开发平台_零代码开发平台_流程设计器_表单引擎_工作流引擎_软件架构

这篇关于基于事件总线EventBus实现邮件推送功能的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1112693

相关文章

C#如何动态创建Label,及动态label事件

《C#如何动态创建Label,及动态label事件》:本文主要介绍C#如何动态创建Label,及动态label事件,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录C#如何动态创建Label,及动态label事件第一点:switch中的生成我们的label事件接着,

使用Sentinel自定义返回和实现区分来源方式

《使用Sentinel自定义返回和实现区分来源方式》:本文主要介绍使用Sentinel自定义返回和实现区分来源方式,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录Sentinel自定义返回和实现区分来源1. 自定义错误返回2. 实现区分来源总结Sentinel自定

Java实现时间与字符串互相转换详解

《Java实现时间与字符串互相转换详解》这篇文章主要为大家详细介绍了Java中实现时间与字符串互相转换的相关方法,文中的示例代码讲解详细,感兴趣的小伙伴可以跟随小编一起学习一下... 目录一、日期格式化为字符串(一)使用预定义格式(二)自定义格式二、字符串解析为日期(一)解析ISO格式字符串(二)解析自定义

opencv图像处理之指纹验证的实现

《opencv图像处理之指纹验证的实现》本文主要介绍了opencv图像处理之指纹验证的实现,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一起学... 目录一、简介二、具体案例实现1. 图像显示函数2. 指纹验证函数3. 主函数4、运行结果三、总结一、

SpringKafka消息发布之KafkaTemplate与事务支持功能

《SpringKafka消息发布之KafkaTemplate与事务支持功能》通过本文介绍的基本用法、序列化选项、事务支持、错误处理和性能优化技术,开发者可以构建高效可靠的Kafka消息发布系统,事务支... 目录引言一、KafkaTemplate基础二、消息序列化三、事务支持机制四、错误处理与重试五、性能优

SpringIntegration消息路由之Router的条件路由与过滤功能

《SpringIntegration消息路由之Router的条件路由与过滤功能》本文详细介绍了Router的基础概念、条件路由实现、基于消息头的路由、动态路由与路由表、消息过滤与选择性路由以及错误处理... 目录引言一、Router基础概念二、条件路由实现三、基于消息头的路由四、动态路由与路由表五、消息过滤

Springboot处理跨域的实现方式(附Demo)

《Springboot处理跨域的实现方式(附Demo)》:本文主要介绍Springboot处理跨域的实现方式(附Demo),具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不... 目录Springboot处理跨域的方式1. 基本知识2. @CrossOrigin3. 全局跨域设置4.

Spring Boot 3.4.3 基于 Spring WebFlux 实现 SSE 功能(代码示例)

《SpringBoot3.4.3基于SpringWebFlux实现SSE功能(代码示例)》SpringBoot3.4.3结合SpringWebFlux实现SSE功能,为实时数据推送提供... 目录1. SSE 简介1.1 什么是 SSE?1.2 SSE 的优点1.3 适用场景2. Spring WebFlu

基于SpringBoot实现文件秒传功能

《基于SpringBoot实现文件秒传功能》在开发Web应用时,文件上传是一个常见需求,然而,当用户需要上传大文件或相同文件多次时,会造成带宽浪费和服务器存储冗余,此时可以使用文件秒传技术通过识别重复... 目录前言文件秒传原理代码实现1. 创建项目基础结构2. 创建上传存储代码3. 创建Result类4.

SpringBoot日志配置SLF4J和Logback的方法实现

《SpringBoot日志配置SLF4J和Logback的方法实现》日志记录是不可或缺的一部分,本文主要介绍了SpringBoot日志配置SLF4J和Logback的方法实现,文中通过示例代码介绍的非... 目录一、前言二、案例一:初识日志三、案例二:使用Lombok输出日志四、案例三:配置Logback一