using System; using System.Collections.Generic; using System.Linq; using System.Text; using System.Threading.Tasks; using YLErp.BLL.Eod; using YLErp.Cache; using YLErp; using YLErp.Helpers; namespace RealTimeCalcPositionService { public class BondCalcKafkTask : IHostedService, IDisposable { private readonly IYcLogger _logger; private KafkaConsumerHelper _rspCalcBondTopicConsumer; private string onRspCalcBondTopic = string.Empty; private string onRspCalcBondTopicGroupId = string.Empty; private readonly CancellationTokenSource _cts = new CancellationTokenSource(); public BondCalcKafkTask() { _logger = LogFactory.GetLogger("BondCalcKafkTask"); onRspCalcBondTopic = Environment.GetEnvironmentVariable("KafkaConfig_OnRspCalcBondTopic"); onRspCalcBondTopicGroupId = Environment.GetEnvironmentVariable("KafkaConfig_OnRspCalcBondTopicGroupId"); _rspCalcBondTopicConsumer = new KafkaConsumerHelper(onRspCalcBondTopicGroupId, onRspCalcBondTopic); } public void Dispose() { _cts.Dispose(); } public Task StartAsync(CancellationToken cancellationToken) { _logger.Info("ClientBalanceTask task is starting."); Task.Run(() => ExecuteTask(_cts.Token), _cts.Token); return Task.CompletedTask; } public async Task StopAsync(CancellationToken cancellationToken) { _logger.Info("ClientBalanceTask task is stopping."); _cts.Cancel(); await Task.CompletedTask; } private void ExecuteTask(CancellationToken stoppingToken) { while (!stoppingToken.IsCancellationRequested) { try { RealtimePnlCalc.ConsumerBondCalcResp(_rspCalcBondTopicConsumer); } catch (Exception ex) { _logger.Error(ex, "收益率接收计算结果服务异常:" + ex.Message); } } } } }