using Org.BouncyCastle.Ocsp; using System.Threading; using YLErp.Abstract; using YLErp.DBModels.Enums; using YLErp.Model.HengTaiModel; using YLErp.Modules.RiskModule; using YLErp.Modules.SwapModule; namespace YLErp.Web.App.KafkaTask { public class ClientReskCheckKafkaTask : IHostedService, IDisposable { private readonly IYcLogger _logger; private readonly CancellationTokenSource _cts = new CancellationTokenSource(); private KafkaConsumerHelper _rspClientRiskCheckConsumer; private string reqClientRiskCheckConsumerTopic = string.Empty; private string reqClientRiskCheckTopicGroupId = string.Empty; private string reqClientRiskCheckTopic = string.Empty; private IKafkaProduce kafkaProduceHelper; public ClientReskCheckKafkaTask(IKafkaProduce kafkaProduce) { _logger = LogFactory.GetLogger("ClientReskCheckKafkaTask"); reqClientRiskCheckTopic = Environment.GetEnvironmentVariable("KafkaConfig_ReqClientRiskCheckTopic"); reqClientRiskCheckTopicGroupId = Environment.GetEnvironmentVariable("KafkaConfig_ReqClientRiskCheckTopicGroupId"); reqClientRiskCheckConsumerTopic = Environment.GetEnvironmentVariable("KafkaConfig_ReqClientRiskCheckConsumerTopic"); _rspClientRiskCheckConsumer = new KafkaConsumerHelper(reqClientRiskCheckTopicGroupId, reqClientRiskCheckTopic); 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 async Task ExecuteTask(CancellationToken cancellationToken) { while (!cancellationToken.IsCancellationRequested) { await Task.Run(()=>{ _rspClientRiskCheckConsumer.Subscribe(msg => { if (!string.IsNullOrEmpty(msg)) { var checkReq = JsonHelper.Deserialize(msg); ClientRiskCheckResp clientRiskCheckResp = new ClientRiskCheckResp(); clientRiskCheckResp.requestId = checkReq.requestId; try { foreach (ClientRiskCheckItemParam req in checkReq.clientRiskCheckItemList) { ClientRiskCheckData clientRiskCheckData = new ClientRiskCheckData(); clientRiskCheckData.orderNo = req.orderNo; clientRiskCheckData.clientRiskCheckItemList = new QuotaMonitorService(new OptUserInfo(0, "客户风控校验", OptUserFrom.Service)).QuotaClientCheck(req); clientRiskCheckData.checkResult = clientRiskCheckData.clientRiskCheckItemList.Count == 0; clientRiskCheckResp.clientRiskCheckDataList.Add(clientRiskCheckData); } } catch (Exception ex) { _logger.Error($"ClientReskCheckKafkaTask Exception ", ex); clientRiskCheckResp.clientRiskCheckDataList.Clear(); foreach (ClientRiskCheckItemParam req in checkReq.clientRiskCheckItemList) { ClientRiskCheckData clientRiskCheckData = new ClientRiskCheckData(); clientRiskCheckData.orderNo = req.orderNo; clientRiskCheckResp.clientRiskCheckDataList.Add(clientRiskCheckData); } clientRiskCheckResp.code = 500; clientRiskCheckResp.msg = ex.Message; } kafkaProduceHelper.Produce(reqClientRiskCheckConsumerTopic, JsonHelper.Serialize(clientRiskCheckResp)); } }); }); } } } }