using Newtonsoft.Json; using YLErp.Abstract; using YLErp.BLL.Eod; using YLErp.Cache; using YLErp.Helpers; using YLErp.Modules.ClientModule; using YLErp.Modules.SwapModule; namespace YLErp.Web.App.KafkaTask { /// /// TRS资金通知消费 /// public class CashNoticeConsumerKafkaTask : IHostedService, IDisposable { private readonly IYcLogger _logger; private readonly CancellationTokenSource _cts = new CancellationTokenSource(); private KafkaConsumerHelper _cashNoticeConsumer; private string cashNoticeTopic = string.Empty; private string cashNoticeGroupId = string.Empty; private IYLCache _yLCache; private IKafkaProduce kafkaProduceHelper; public CashNoticeConsumerKafkaTask(IYLCache yLCache, IKafkaProduce kafkaProduce) { _logger = LogFactory.GetLogger("CashNoticeConsumerKafkaTask"); cashNoticeTopic = Environment.GetEnvironmentVariable("KafkaConfig_YiLian_CashNoticeConsumerTopic"); cashNoticeGroupId = Environment.GetEnvironmentVariable("KafkaConfig_YiLian_CashNoticeConsumerGroup"); _cashNoticeConsumer = new KafkaConsumerHelper(cashNoticeGroupId, cashNoticeTopic); _yLCache = yLCache; kafkaProduceHelper = kafkaProduce; } public void Dispose() { _cts.Dispose(); } public Task StartAsync(CancellationToken cancellationToken) { Task.Run(() => ExecuteTask(_cts.Token), _cts.Token); return Task.CompletedTask; } public async Task StopAsync(CancellationToken cancellationToken) { _cts.Cancel(); await Task.CompletedTask; } private void ExecuteTask(CancellationToken cancellationToken) { while (!cancellationToken.IsCancellationRequested) { try { var service = new SwapConsumerService(OptUserInfo.SystemUser); service.SetKafKaProduce(kafkaProduceHelper); } catch (Exception ex) { _logger.Error("CashNoticeConsumerKafkaTask Exception", ex); } } } } }