using Newtonsoft.Json; using YLErp; using YLErp.Abstract; using YLErp.BLL; using YLErp.BLL.Eod; using YLErp.Cache; using YLErp.Helpers; using YLErp.Model; namespace RealTimeCalcPositionService { public class ClientNoDMABalanceTask : IHostedService, IDisposable { private readonly IYcLogger _logger; private readonly CancellationTokenSource _cts = new CancellationTokenSource(); private IKafkaProduce kafkaProduceHelper; private string onRspAccountCapitalTopicTopic = string.Empty; private Dictionary clientDic = new Dictionary(); private IYLCache _yLCache; public ClientNoDMABalanceTask(IKafkaProduce kafkaProduce, IYLCache yLCache) { _logger = LogFactory.GetLogger("ClientNoDMABalanceTask"); onRspAccountCapitalTopicTopic = Environment.GetEnvironmentVariable("KafkaConfig_OnRspAccountCapitalTopic"); kafkaProduceHelper = kafkaProduce; _yLCache = yLCache; } public void Dispose() { _cts.Dispose(); } public Task StartAsync(CancellationToken cancellationToken) { _logger.Info("ClientNoDMABalanceTask task is starting."); Task.Run(() => ExecuteTask(_cts.Token), _cts.Token); return Task.CompletedTask; } private void ExecuteTask(CancellationToken stoppingToken) { while (!stoppingToken.IsCancellationRequested) { try { //系统交易日 var valuedate = valuedateBLL.ValueDate; var clientSettles = new RealTimeClientBanlanceService(new OptUserInfo(0, "实时客户资金服务", OptUserFrom.Service)).GetBalances(); foreach (var cb in clientSettles) { var AvailableMoney = cb.AvailableAmount + cb.FrozenMarginMoney; Result result = new Result(); try { var obj = new ClientBalanceForTrsResponse { AvailableAmount = Math.Round(cb.AvailableAmount, 2), PositionPv = cb.RoundedPositionPv, PositionPnl = cb.RoundedPositionPnl, ClientId = cb.ClientId, ClientType = cb.ClientType, Credit = cb.TotalCredit, AvailableMoney = AvailableMoney }; result.success = true; result.obj = obj; if (_yLCache != null) { _yLCache.StringSet("ClientBalance:" + cb.ClientId, obj); _yLCache.HashSet("risk:cash:balance:amount", cb.ClientId.ToString(), (decimal)obj.AvailableMoney); } } catch (Exception ex) { result.msg = ex.Message; result.success = false; } string resultStr = JsonConvert.SerializeObject(result); string newEncryStr = DataProtectHelper.Encrypt(resultStr); bool needProduce = false; if (clientDic.TryGetValue(cb.ClientId, out string encryStr)) { if (encryStr != newEncryStr) { needProduce = true; clientDic[cb.ClientId] = newEncryStr; } } else { clientDic.Add(cb.ClientId, newEncryStr); needProduce = true; } if (needProduce) { // 使用原有的topic而不是新增的RealtimeCalcAccountBalance kafkaProduceHelper.Produce(onRspAccountCapitalTopicTopic, resultStr); } } Thread.Sleep(3000); } catch (Exception ex) { _logger.Error(ex, "普通实时客户资金服务异常:" + ex.Message); } } } public async Task StopAsync(CancellationToken cancellationToken) { _logger.Info("ClientNoDMABalanceTask task is stopping."); _cts.Cancel(); await Task.CompletedTask; } } }