1. 新增数据库唯一约束和Message_Outbox表脚本 2. 新增雪花ID、Hangfire存储、MQ重试等配置实体 3. 重构各项目雪花ID生成逻辑,改为从配置读取WorkerId 4. 优化积分服务分页查询、用户背包更新逻辑 5. 新增JWT令牌Redis过期刷新逻辑 6. 完善RabbitMQ死信队列消息头信息 7. 新增可靠MQ消息发布服务和Outbox派发后台服务 8. 替换原有RabbitMQ直接发送为Outbox可靠发布 9. 优化签到服务逻辑,新增重复签到校验和补签卡扣减逻辑 10. 修复自动铺码消费逻辑,新增点阵页预占和释放机制
113 lines
4.5 KiB
C#
113 lines
4.5 KiB
C#
using Microsoft.Extensions.Logging;
|
||
using QYZH.InteractiveMagazine.Infrastructure.RabbitMQ;
|
||
using QYZH.InteractiveMagazine.IService;
|
||
using QYZH.InteractiveMagazine.Models.Common;
|
||
using QYZH.InteractiveMagazine.Models.Dto;
|
||
using QYZH.InteractiveMagazine.Models.Dto.RabbitMQ;
|
||
using QYZH.InteractiveMagazine.Models.Entity;
|
||
using QYZH.InteractiveMagazine.Models.Enum;
|
||
using QYZH.InteractiveMagazine.Repository;
|
||
using System.Text.Json;
|
||
|
||
namespace QYZH.InteractiveMagazine.Service;
|
||
|
||
/// <summary>
|
||
/// MQ消息发布服务。
|
||
/// </summary>
|
||
public class MessagePublishService(
|
||
IRabbitMQService rabbitMQService,
|
||
ILogger<MessagePublishService> logger) : BaseRepository<MessageOutbox>, IMessagePublishService
|
||
{
|
||
/// <summary>
|
||
/// 可靠发布消息,默认写入Outbox。
|
||
/// </summary>
|
||
public async Task<MessagePublishResult> PublishAsync<T>(MessagePublishInput<T> input, CancellationToken cancellationToken = default)
|
||
{
|
||
ValidateInput(input);
|
||
var message = BuildOutboxMessage(input);
|
||
await Context.Insertable(message).ExecuteCommandAsync(cancellationToken);
|
||
return new MessagePublishResult
|
||
{
|
||
Success = true,
|
||
OutboxIds = [message.Id],
|
||
Message = "消息已写入Outbox"
|
||
};
|
||
}
|
||
|
||
/// <summary>
|
||
/// 批量可靠发布消息,默认写入Outbox。
|
||
/// </summary>
|
||
public async Task<MessagePublishResult> PublishBatchAsync<T>(IEnumerable<MessagePublishInput<T>> inputs, CancellationToken cancellationToken = default)
|
||
{
|
||
var inputList = inputs.ToList();
|
||
BusinessException.ThrowIf(inputList.Count == 0, "消息发布列表不能为空", ResultCode.BAD_REQUEST);
|
||
inputList.ForEach(ValidateInput);
|
||
|
||
var messages = inputList.Select(BuildOutboxMessage).ToList();
|
||
await Context.Insertable(messages).ExecuteCommandAsync(cancellationToken);
|
||
return new MessagePublishResult
|
||
{
|
||
Success = true,
|
||
OutboxIds = messages.Select(x => x.Id).ToList(),
|
||
Message = "消息已批量写入Outbox"
|
||
};
|
||
}
|
||
|
||
/// <summary>
|
||
/// 直接发布消息,不写Outbox。
|
||
/// </summary>
|
||
public async Task<MessagePublishResult> PublishDirectAsync<T>(MessagePublishInput<T> input, CancellationToken cancellationToken = default)
|
||
{
|
||
ValidateInput(input);
|
||
var sent = await rabbitMQService.SendAsync(new RabbitMQSendParam
|
||
{
|
||
Exchange = input.Exchange,
|
||
Queue = input.Queue,
|
||
RoutingKey = input.RoutingKey,
|
||
Data = input.Data!
|
||
}, cancellationToken);
|
||
|
||
if (!sent)
|
||
{
|
||
logger.LogError("MQ直接发布失败,Exchange: {Exchange}, Queue: {Queue}, RoutingKey: {RoutingKey}, BusinessType: {BusinessType}, BusinessId: {BusinessId}",
|
||
input.Exchange, input.Queue, input.RoutingKey, input.BusinessType, input.BusinessId);
|
||
throw new BusinessException("MQ消息发送失败", ResultCode.GLOBAL_ERROR);
|
||
}
|
||
|
||
return new MessagePublishResult
|
||
{
|
||
Success = true,
|
||
Message = "消息已直接发送"
|
||
};
|
||
}
|
||
|
||
private static MessageOutbox BuildOutboxMessage<T>(MessagePublishInput<T> input)
|
||
{
|
||
var now = DateTime.Now;
|
||
return new MessageOutbox
|
||
{
|
||
Exchange = input.Exchange,
|
||
Queue = input.Queue,
|
||
RoutingKey = input.RoutingKey,
|
||
Payload = JsonSerializer.Serialize(input.Data),
|
||
BusinessType = input.BusinessType,
|
||
BusinessId = input.BusinessId,
|
||
Status = (int)MessageOutboxStatusEnum.Pending,
|
||
CreatedBy = "System",
|
||
CreatedAt = now,
|
||
UpdatedBy = "System",
|
||
UpdatedAt = now
|
||
};
|
||
}
|
||
|
||
private static void ValidateInput<T>(MessagePublishInput<T> input)
|
||
{
|
||
BusinessException.ThrowIf(string.IsNullOrWhiteSpace(input.Exchange), "MQ交换机不能为空", ResultCode.BAD_REQUEST);
|
||
BusinessException.ThrowIf(string.IsNullOrWhiteSpace(input.Queue), "MQ队列不能为空", ResultCode.BAD_REQUEST);
|
||
BusinessException.ThrowIf(string.IsNullOrWhiteSpace(input.RoutingKey), "MQ路由键不能为空", ResultCode.BAD_REQUEST);
|
||
BusinessException.ThrowIf(string.IsNullOrWhiteSpace(input.BusinessType), "MQ业务类型不能为空", ResultCode.BAD_REQUEST);
|
||
BusinessException.ThrowIf(input.BusinessId <= 0, "MQ业务ID无效", ResultCode.BAD_REQUEST);
|
||
BusinessException.ThrowIf(input.Data == null, "MQ消息数据不能为空", ResultCode.BAD_REQUEST);
|
||
}
|
||
}
|