2026-06-29 16:34:26 +08:00
|
|
|
using Newtonsoft.Json;
|
|
|
|
|
using QYZH.InteractiveMagazine.Infrastructure.OSS;
|
|
|
|
|
using SqlSugar;
|
2026-06-11 18:01:07 +08:00
|
|
|
using System.Text;
|
2026-06-29 16:34:26 +08:00
|
|
|
using System.Threading.Channels;
|
|
|
|
|
using Yitter.IdGenerator;
|
2026-06-11 18:01:07 +08:00
|
|
|
|
|
|
|
|
namespace QYZH.InteractiveMagazine.WorkService.Consumers;
|
|
|
|
|
|
|
|
|
|
/// <summary>
|
|
|
|
|
/// 期刊任务接收消费者(示例)
|
|
|
|
|
/// </summary>
|
2026-06-29 16:34:26 +08:00
|
|
|
public class JournalTaskReceiveConsumer(ILogger<JournalTaskReceiveConsumer> logger, IConfiguration configuration,
|
|
|
|
|
IServiceScopeFactory scopeFactory,
|
|
|
|
|
IWebHostEnvironment webHostEnvironment,
|
|
|
|
|
IHttpClientFactory httpClientFactory,
|
|
|
|
|
OssService ossService) : IQueueConsumer
|
2026-06-11 18:01:07 +08:00
|
|
|
{
|
|
|
|
|
|
2026-06-22 16:56:05 +08:00
|
|
|
public string Exchange => "ex.journal";
|
|
|
|
|
|
|
|
|
|
public string QueueName => "mq.journal.task.receive";
|
|
|
|
|
|
|
|
|
|
public string RoutingKey => "rk.journal.task.receive";
|
2026-06-11 18:01:07 +08:00
|
|
|
|
2026-06-29 16:34:26 +08:00
|
|
|
public async Task HandleAsync(byte[] body, CancellationToken cancellationToken = default)
|
2026-06-11 18:01:07 +08:00
|
|
|
{
|
2026-06-29 16:34:26 +08:00
|
|
|
var message = Encoding.UTF8.GetString(body);
|
|
|
|
|
logger.LogInformation("收到期刊任务消息: {Message}", message);
|
2026-06-11 18:01:07 +08:00
|
|
|
|
2026-06-29 16:34:26 +08:00
|
|
|
using var scope = scopeFactory.CreateScope();
|
|
|
|
|
var client = scope.ServiceProvider.GetRequiredService<ISqlSugarClient>();
|
2026-06-11 18:01:07 +08:00
|
|
|
// TODO: 在此编写具体的消息处理逻辑
|
2026-06-29 16:34:26 +08:00
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
var data = System.Text.Json.JsonSerializer.Deserialize<QuestionData>(message);
|
|
|
|
|
|
2026-06-11 18:01:07 +08:00
|
|
|
|
2026-06-29 16:34:26 +08:00
|
|
|
|
|
|
|
|
client.Ado.CommitTran();
|
|
|
|
|
}
|
|
|
|
|
catch (Exception ex)
|
|
|
|
|
{
|
|
|
|
|
client.Ado.RollbackTran();
|
|
|
|
|
logger.LogError(ex.Message + ex.StackTrace);
|
|
|
|
|
}
|
2026-06-11 18:01:07 +08:00
|
|
|
await Task.CompletedTask;
|
|
|
|
|
}
|
|
|
|
|
|
2026-06-29 16:34:26 +08:00
|
|
|
public Task OnErrorAsync(byte[] body, Exception exception)
|
2026-06-11 18:01:07 +08:00
|
|
|
{
|
2026-06-29 16:34:26 +08:00
|
|
|
var message = Encoding.UTF8.GetString(body);
|
|
|
|
|
logger.LogError(exception, "处理期刊任务消息失败: {Message}", message);
|
2026-06-11 18:01:07 +08:00
|
|
|
return Task.CompletedTask;
|
|
|
|
|
}
|
2026-06-29 16:34:26 +08:00
|
|
|
|
|
|
|
|
|
2026-06-11 18:01:07 +08:00
|
|
|
}
|
2026-06-29 16:34:26 +08:00
|
|
|
|
|
|
|
|
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; }
|
|
|
|
|
}
|