Files
QYZH.InteractiveMagazine/QYZH.InteractiveMagazine.WorkService/Consumers/JournalTaskReceiveConsumer.cs

85 lines
2.5 KiB
C#
Raw Normal View History

using Newtonsoft.Json;
using QYZH.InteractiveMagazine.Infrastructure.OSS;
using SqlSugar;
using System.Text;
using System.Threading.Channels;
using Yitter.IdGenerator;
namespace QYZH.InteractiveMagazine.WorkService.Consumers;
/// <summary>
/// 期刊任务接收消费者(示例)
/// </summary>
public class JournalTaskReceiveConsumer(ILogger<JournalTaskReceiveConsumer> logger, IConfiguration configuration,
IServiceScopeFactory scopeFactory,
IWebHostEnvironment webHostEnvironment,
IHttpClientFactory httpClientFactory,
OssService ossService) : IQueueConsumer
{
public string Exchange => "ex.journal";
public string QueueName => "mq.journal.task.receive";
public string RoutingKey => "rk.journal.task.receive";
public async Task HandleAsync(byte[] body, CancellationToken cancellationToken = default)
{
var message = Encoding.UTF8.GetString(body);
logger.LogInformation("收到期刊任务消息: {Message}", message);
using var scope = scopeFactory.CreateScope();
var client = scope.ServiceProvider.GetRequiredService<ISqlSugarClient>();
// TODO: 在此编写具体的消息处理逻辑
try
{
var data = System.Text.Json.JsonSerializer.Deserialize<QuestionData>(message);
client.Ado.CommitTran();
}
catch (Exception ex)
{
client.Ado.RollbackTran();
logger.LogError(ex.Message + ex.StackTrace);
}
await Task.CompletedTask;
}
public Task OnErrorAsync(byte[] body, Exception exception)
{
var message = Encoding.UTF8.GetString(body);
logger.LogError(exception, "处理期刊任务消息失败: {Message}", message);
return Task.CompletedTask;
}
}
public class QuestionData
{
public long UserId { get; set; }
public long JournalId { get; set; }
public long PageId { get; set; }
public string PageAnswerUrl { get; set; }
public Question[] Questions { get; set; }
public DateTime CreatedTime { get; set; }
}
public class Question
{
public long Id { get; set; }
public string Url { get; set; }
public string[] AnswerUrl { get; set; }
public DateTime AnswerStartTime { get; set; }
public DateTime AnswerEndTime { get; set; }
public int AnswerTime { get; set; }
public int BreakCount { get; set; }
public List<BreakTime> BreakTimes { get; set; }
}
public class BreakTime
{
public DateTime Time { get; set; }
public long WaitTime { get; set; }
}