using Dpz.Core.Entity.Base;
using Dpz.Core.MessageQueue.Abstractions;
using Dpz.Core.MessageQueue.Configuration;
using Dpz.Core.MessageQueue.Models;
using Dpz.Core.MessageQueue.RabbitMQ;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
namespace Dpz.Core.MessageQueue.Extensions;
/// <summary>
/// 消息队列服务扩展方法
/// </summary>
public static class MessageQueueServiceExtensions
{
/// <summary>
/// 添加RabbitMQ服务
/// </summary>
/// <param name="services">服务集合</param>
/// <param name="configuration">配置</param>
/// <param name="configSection">配置节名称,默认为"RabbitMQ"</param>
/// <returns></returns>
public static IServiceCollection AddRabbitMQ(
this IServiceCollection services,
IConfiguration configuration,
string configSection = RabbitMQOptions.SectionName
)
{
// 注册配置
var section = configuration.GetSection(configSection);
services.Configure<RabbitMQOptions>(section);
// 注册连接工厂(单例)
services.AddSingleton<IRabbitMQConnectionFactory, RabbitMQConnectionFactory>();
// 注册路由约定(单例)
services.AddSingleton<IMessageRoutingConvention, DefaultMessageRoutingConvention>();
// 注册发布者(单例)
services.AddSingleton(typeof(IMessagePublisher<>), typeof(RabbitMQPublisher<>));
// 注册 Outbox 空实现(默认不启用,可通过 AddMessageOutbox() 替换为 MongoDB 实现)
services.AddSingleton<IMessageOutboxStore, NullMessageOutboxStore>();
services.AddSingleton<IMessageOutboxRetryPublisher, NullMessageOutboxRetryPublisher>();
return services;
}
/// <summary>
/// 添加消息消费者
/// 每个消息类型配置一个消费者和处理器
/// </summary>
/// <typeparam name="TMessage">消息类型</typeparam>
/// <typeparam name="THandler">处理器类型</typeparam>
/// <param name="services">服务集合</param>
/// <returns></returns>
public static IServiceCollection AddMessageConsumer<TMessage, THandler>(
this IServiceCollection services
)
where TMessage : MessageBase
where THandler : class, IMessageHandler<TMessage>
{
// 注册处理器(Scoped,每次处理消息创建新实例)
services.AddScoped<IMessageHandler<TMessage>, THandler>();
// 注册消费者后台服务
services.AddHostedService<RabbitMQConsumerBackgroundService<TMessage>>();
return services;
}
/// <summary>
/// 添加消息消费者(支持返回结果)
/// 每个消息类型配置一个消费者和处理器
/// </summary>
/// <typeparam name="TMessage">消息类型</typeparam>
/// <typeparam name="THandler">处理器类型</typeparam>
/// <typeparam name="TResult">返回结果类型</typeparam>
/// <param name="services">服务集合</param>
/// <returns></returns>
public static IServiceCollection AddMessageConsumer<TMessage, THandler, TResult>(
this IServiceCollection services
)
where TMessage : MessageBase
where TResult : MessageHandlerResult
where THandler : class, IMessageHandler<TMessage, TResult>
{
// 注册处理器(Scoped,每次处理消息创建新实例)
services.AddScoped<IMessageHandler<TMessage, TResult>, THandler>();
// 注册消费者后台服务(支持返回结果)
services.AddHostedService<RabbitMQConsumerBackgroundServiceWithResult<TMessage, TResult>>();
return services;
}
/// <summary>
/// 添加自定义路由约定
/// </summary>
/// <typeparam name="TConvention">自定义路由约定类型</typeparam>
/// <param name="services">服务集合</param>
/// <returns></returns>
public static IServiceCollection AddCustomRoutingConvention<TConvention>(
this IServiceCollection services
)
where TConvention : class, IMessageRoutingConvention
{
services.AddSingleton<IMessageRoutingConvention, TConvention>();
return services;
}
/// <summary>
/// 添加批次追踪服务
/// (多实例部署时需要配置分布式缓存如Redis)
/// </summary>
/// <param name="services">服务集合</param>
/// <returns></returns>
public static IServiceCollection AddBatchTracking(this IServiceCollection services)
{
services.AddSingleton<IBatchTracker, DistributedBatchTracker>();
return services;
}
/// <summary>
/// 添加 Outbox 后台补发服务。
/// </summary>
/// <param name="services">服务集合。</param>
/// <returns>注册后的服务集合。</returns>
public static IServiceCollection AddMessageOutboxRetryWorker(this IServiceCollection services)
{
services.AddSingleton<IMessageOutboxRetryService, MessageOutboxRetryService>();
services.AddHostedService<MessageOutboxRetryBackgroundService>();
return services;
}
}
⚠⚠ 以下内容为AI分析的结果,请根据实际情况进行判断。
这是一个消息队列服务的依赖注入扩展类,主要用于在.NET应用程序中配置和注册RabbitMQ相关的服务。下面我来详细解释这个代码:
类概述
MessageQueueServiceExtensions 是一个静态扩展类,为 IServiceCollection 提供了一系列用于配置消息队列功能的扩展方法。
主要方法解析
1. AddRabbitMQ 方法
public static IServiceCollection AddRabbitMQ(
this IServiceCollection services,
IConfiguration configuration,
string configSection = RabbitMQOptions.SectionName
)
作用: 注册RabbitMQ的基础服务 注册的服务:
RabbitMQOptions配置选项IRabbitMQConnectionFactory连接工厂(单例)IMessageRoutingConvention路由约定(单例)IMessagePublisher<>消息发布者(单例)IMessageOutboxStore和IMessageOutboxRetryPublisher的空实现(默认不启用Outbox模式)
2. AddMessageConsumer<TMessage, THandler> 方法
public static IServiceCollection AddMessageConsumer<TMessage, THandler>(
this IServiceCollection services
)
作用: 添加消息消费者和处理器 泛型约束:
TMessage必须继承自MessageBaseTHandler必须实现IMessageHandler<TMessage>
注册的服务:
- 消息处理器(Scoped生命周期)
- 后台消费者服务
3. AddMessageConsumer<TMessage, THandler, TResult> 方法
作用: 添加支持返回结果的消息消费者 额外约束: TResult 必须继承自 MessageHandlerResult 特点: 使用带返回结果的后台服务实现
4. AddCustomRoutingConvention<TConvention> 方法
作用: 允许用户注册自定义的消息路由约定 用途: 定制消息的交换机、队列、路由键等命名规则
5. AddBatchTracking 方法
作用: 添加批次追踪服务 注意: 多实例部署时需要配置分布式缓存(如Redis)来共享批次状态
6. AddMessageOutboxRetryWorker 方法
作用: 添加Outbox模式的重试补发服务 注册的服务:
IMessageOutboxRetryService重试服务MessageOutboxRetryBackgroundService后台重试工作服务
设计模式和特点
- 扩展方法模式: 为
IServiceCollection提供流畅的API - 依赖注入: 遵循.NET Core的DI容器模式
- 泛型约束: 确保类型安全
- 生命周期管理:
- 单例:连接工厂、发布者等共享资源
- Scoped:消息处理器,每次处理创建新实例
- 可扩展性: 支持自定义路由约定和处理器
- Outbox模式: 支持事务性消息发送的可靠性保证
使用场景
这个扩展类适用于需要在.NET应用中集成RabbitMQ消息队列的场景,提供了完整的消息发布、消费、路由和可靠性保证的功能。
AI 正在分析代码…
评论加载中...