using YLErp.Abstract; using YLErp.Helpers; using YLErp.Modules.SwapModule; namespace YLErp.Web.App.KafkaTask { /// /// 交易对手方信息推送回应衡泰Kafka /// public class ClientInfoConsumerKafkaTask : IHostedService, IDisposable { private readonly IYcLogger _logger; private readonly CancellationTokenSource _cts = new CancellationTokenSource(); private string onRspCounterPartyInsertTopic = string.Empty; private string onRspCounterPartyInsertTopicGroupId = string.Empty; private KafkaConsumerHelper _rspCounterPartyInsertConsumer; private IKafkaProduce kafkaProduceHelper; public ClientInfoConsumerKafkaTask(IKafkaProduce kafkaProduce) { _logger = LogFactory.GetLogger("ClientInfoConsumerKafkaTask"); onRspCounterPartyInsertTopic = Environment.GetEnvironmentVariable("KafkaConfig_OnRspCounterPartyInsertTopic"); onRspCounterPartyInsertTopicGroupId = Environment.GetEnvironmentVariable("KafkaConfig_OnRspCounterPartyInsertTopicGroupId"); _rspCounterPartyInsertConsumer = new KafkaConsumerHelper(onRspCounterPartyInsertTopicGroupId, onRspCounterPartyInsertTopic); kafkaProduceHelper = kafkaProduce; } public void Dispose() { _cts.Dispose(); } public Task StartAsync(CancellationToken cancellationToken) { _logger.Info("ConsumerClient started"); Task.Run(() => ExecuteTask(_cts.Token), _cts.Token); return Task.CompletedTask; } public async Task StopAsync(CancellationToken cancellationToken) { _logger.Info("ConsumerClient stoped"); _cts.Cancel(); await Task.CompletedTask; } private async Task ExecuteTask(CancellationToken cancellationToken) { while (!cancellationToken.IsCancellationRequested) { try { var pushService = new SwapPushService(OptUserInfo.SystemUser); pushService.SetKafKaProduce(kafkaProduceHelper); // 执行任务 await Task.Run(() => { new SwapPushService(OptUserInfo.SystemUser).ConsumerClientResp(_rspCounterPartyInsertConsumer); }, cancellationToken); } catch (Exception ex) { _logger.Error("ConsumerClient Exception", ex); } } } } }