基于事件总线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

相关文章

python使用fastapi实现多语言国际化的操作指南

《python使用fastapi实现多语言国际化的操作指南》本文介绍了使用Python和FastAPI实现多语言国际化的操作指南,包括多语言架构技术栈、翻译管理、前端本地化、语言切换机制以及常见陷阱和... 目录多语言国际化实现指南项目多语言架构技术栈目录结构翻译工作流1. 翻译数据存储2. 翻译生成脚本

如何通过Python实现一个消息队列

《如何通过Python实现一个消息队列》这篇文章主要为大家详细介绍了如何通过Python实现一个简单的消息队列,文中的示例代码讲解详细,感兴趣的小伙伴可以跟随小编一起学习一下... 目录如何通过 python 实现消息队列如何把 http 请求放在队列中执行1. 使用 queue.Queue 和 reque

Python如何实现PDF隐私信息检测

《Python如何实现PDF隐私信息检测》随着越来越多的个人信息以电子形式存储和传输,确保这些信息的安全至关重要,本文将介绍如何使用Python检测PDF文件中的隐私信息,需要的可以参考下... 目录项目背景技术栈代码解析功能说明运行结php果在当今,数据隐私保护变得尤为重要。随着越来越多的个人信息以电子形

使用 sql-research-assistant进行 SQL 数据库研究的实战指南(代码实现演示)

《使用sql-research-assistant进行SQL数据库研究的实战指南(代码实现演示)》本文介绍了sql-research-assistant工具,该工具基于LangChain框架,集... 目录技术背景介绍核心原理解析代码实现演示安装和配置项目集成LangSmith 配置(可选)启动服务应用场景

使用Python快速实现链接转word文档

《使用Python快速实现链接转word文档》这篇文章主要为大家详细介绍了如何使用Python快速实现链接转word文档功能,文中的示例代码讲解详细,感兴趣的小伙伴可以跟随小编一起学习一下... 演示代码展示from newspaper import Articlefrom docx import

前端原生js实现拖拽排课效果实例

《前端原生js实现拖拽排课效果实例》:本文主要介绍如何实现一个简单的课程表拖拽功能,通过HTML、CSS和JavaScript的配合,我们实现了课程项的拖拽、放置和显示功能,文中通过实例代码介绍的... 目录1. 效果展示2. 效果分析2.1 关键点2.2 实现方法3. 代码实现3.1 html部分3.2

Java深度学习库DJL实现Python的NumPy方式

《Java深度学习库DJL实现Python的NumPy方式》本文介绍了DJL库的背景和基本功能,包括NDArray的创建、数学运算、数据获取和设置等,同时,还展示了如何使用NDArray进行数据预处理... 目录1 NDArray 的背景介绍1.1 架构2 JavaDJL使用2.1 安装DJL2.2 基本操

最长公共子序列问题的深度分析与Java实现方式

《最长公共子序列问题的深度分析与Java实现方式》本文详细介绍了最长公共子序列(LCS)问题,包括其概念、暴力解法、动态规划解法,并提供了Java代码实现,暴力解法虽然简单,但在大数据处理中效率较低,... 目录最长公共子序列问题概述问题理解与示例分析暴力解法思路与示例代码动态规划解法DP 表的构建与意义动

java父子线程之间实现共享传递数据

《java父子线程之间实现共享传递数据》本文介绍了Java中父子线程间共享传递数据的几种方法,包括ThreadLocal变量、并发集合和内存队列或消息队列,并提醒注意并发安全问题... 目录通过 ThreadLocal 变量共享数据通过并发集合共享数据通过内存队列或消息队列共享数据注意并发安全问题总结在 J

SpringBoot+MyBatis-Flex配置ProxySQL的实现步骤

《SpringBoot+MyBatis-Flex配置ProxySQL的实现步骤》本文主要介绍了SpringBoot+MyBatis-Flex配置ProxySQL的实现步骤,文中通过示例代码介绍的非常详... 目录 目标 步骤 1:确保 ProxySQL 和 mysql 主从同步已正确配置ProxySQL 的