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<> 消息发布者(单例)
  • IMessageOutboxStoreIMessageOutboxRetryPublisher 的空实现(默认不启用Outbox模式)

2. AddMessageConsumer<TMessage, THandler> 方法

public static IServiceCollection AddMessageConsumer<TMessage, THandler>(
    this IServiceCollection services
)

作用: 添加消息消费者和处理器 泛型约束:

  • TMessage 必须继承自 MessageBase
  • THandler 必须实现 IMessageHandler<TMessage>

注册的服务:

  • 消息处理器(Scoped生命周期)
  • 后台消费者服务

3. AddMessageConsumer<TMessage, THandler, TResult> 方法

作用: 添加支持返回结果的消息消费者 额外约束: TResult 必须继承自 MessageHandlerResult 特点: 使用带返回结果的后台服务实现

4. AddCustomRoutingConvention<TConvention> 方法

作用: 允许用户注册自定义的消息路由约定 用途: 定制消息的交换机、队列、路由键等命名规则

5. AddBatchTracking 方法

作用: 添加批次追踪服务 注意: 多实例部署时需要配置分布式缓存(如Redis)来共享批次状态

6. AddMessageOutboxRetryWorker 方法

作用: 添加Outbox模式的重试补发服务 注册的服务:

  • IMessageOutboxRetryService 重试服务
  • MessageOutboxRetryBackgroundService 后台重试工作服务

设计模式和特点

  1. 扩展方法模式: 为 IServiceCollection 提供流畅的API
  2. 依赖注入: 遵循.NET Core的DI容器模式
  3. 泛型约束: 确保类型安全
  4. 生命周期管理:
    • 单例:连接工厂、发布者等共享资源
    • Scoped:消息处理器,每次处理创建新实例
  5. 可扩展性: 支持自定义路由约定和处理器
  6. Outbox模式: 支持事务性消息发送的可靠性保证

使用场景

这个扩展类适用于需要在.NET应用中集成RabbitMQ消息队列的场景,提供了完整的消息发布、消费、路由和可靠性保证的功能。

评论加载中...