using Newtonsoft.Json; using YLErp.Abstract; using YLErp.Modules.SwapModule; namespace YLErp.Web.App.KafkaTask { /// /// trs交易端单个客户资金监控 /// public class ClientBalanceMonitorSingleKafkaTask : IHostedService, IDisposable { private readonly IYcLogger _logger; private readonly CancellationTokenSource _cts = new CancellationTokenSource(); private string onRspClientMonitorSingleConsumerTopic = string.Empty; private string reqClientMonitorSingleTopicGroupId = string.Empty; private string reqClientMonitorSingleTopic = string.Empty; private KafkaConsumerHelper _kafkaConsumer; private IKafkaProduce kafkaProduceHelper; public ClientBalanceMonitorSingleKafkaTask(IKafkaProduce kafkaProduce) { _logger = LogFactory.GetLogger("ClientBalanceMonitorSingleKafkaTask"); onRspClientMonitorSingleConsumerTopic = Environment.GetEnvironmentVariable("KafkaConfig_OnRspClientMonitorSingleConsumerTopic"); reqClientMonitorSingleTopicGroupId = Environment.GetEnvironmentVariable("KafkaConfig_ReqClientMonitorSingleTopicGroupId"); reqClientMonitorSingleTopic = Environment.GetEnvironmentVariable("KafkaConfig_ReqClientMonitorSingleTopic"); _kafkaConsumer = new KafkaConsumerHelper(reqClientMonitorSingleTopicGroupId, reqClientMonitorSingleTopic); kafkaProduceHelper = kafkaProduce; } public void Dispose() { _cts.Dispose(); } public Task StartAsync(CancellationToken cancellationToken) { _logger.Info("ClientBalanceMonitorSingleKafkaTask started"); Task.Run(() => ExecuteTask(_cts.Token), _cts.Token); return Task.CompletedTask; } public async Task StopAsync(CancellationToken cancellationToken) { _logger.Info("ClientBalanceMonitorSingleKafkaTask stoped"); _cts.Cancel(); await Task.CompletedTask; } private async Task ExecuteTask(CancellationToken cancellationToken) { while (!cancellationToken.IsCancellationRequested) { try { _kafkaConsumer.Subscribe(async msg => { if (!string.IsNullOrEmpty(msg)) { ClientBalanceMonitorFroTrsRequest req = JsonConvert.DeserializeObject(msg); var swapMonitorService = new SwapMonitorService(OptUserInfo.SystemUser); // 执行任务 await Task.Run(() => { Result result = new Result(); var res = new TrsRespone() { Req= req }; try { res.Res= swapMonitorService.GetMonitorForTrsBuyDailyRespone(req); result.obj = res; result.success = true; } catch (Exception ex) { result.msg = ex.Message; result.success = false; } kafkaProduceHelper.Produce(onRspClientMonitorSingleConsumerTopic, JsonConvert.SerializeObject(result)); }, cancellationToken); } }); } catch (Exception ex) { _logger.Error("ClientBalanceMonitorSingleKafkaTask Exception", ex); } } } } }