You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

TDengineService.cs 59KB

3 months ago
3 months ago
1 year ago
1 year ago
1 year ago
11 months ago
11 months ago
11 months ago
1 year ago
1 year ago
1 year ago
1 year ago
11 months ago
4 days ago
4 days ago
4 days ago
4 days ago
4 days ago
4 days ago
4 months ago
4 months ago
3 months ago
3 months ago
3 months ago
3 months ago
3 months ago
3 months ago
3 months ago
3 months ago
4 months ago
4 months ago
4 months ago
3 months ago
4 months ago
4 months ago
4 months ago
3 months ago
3 months ago
3 months ago
4 months ago
3 months ago
4 months ago
4 months ago
4 months ago
4 months ago
4 months ago
4 months ago
4 months ago
4 months ago
4 months ago
4 months ago
4 months ago
4 months ago
4 months ago
4 months ago
4 months ago
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396
  1. using HealthMonitor.Common;
  2. using HealthMonitor.Common.helper;
  3. using HealthMonitor.Model.Config;
  4. using HealthMonitor.Model.Service.Mapper;
  5. using HealthMonitor.Service.Biz.db.Dto;
  6. using HealthMonitor.Util.Models;
  7. using Microsoft.EntityFrameworkCore.Metadata.Internal;
  8. using Microsoft.Extensions.Logging;
  9. using Microsoft.Extensions.Options;
  10. using Newtonsoft.Json;
  11. using Newtonsoft.Json.Linq;
  12. using SqlSugar;
  13. using SqlSugar.DbConvert;
  14. using SqlSugar.TDengine;
  15. using System;
  16. using System.Collections.Generic;
  17. using System.Data;
  18. using System.Linq;
  19. using System.Linq.Expressions;
  20. using System.Reflection;
  21. using System.Text;
  22. using System.Threading.Tasks;
  23. using System.Xml.Linq;
  24. using TDengineDriver;
  25. using TDengineDriver.Impl;
  26. using TDengineTMQ;
  27. using HealthMonitor.Service.Cache;
  28. using System.Text.RegularExpressions;
  29. using Etcdserverpb;
  30. using static Microsoft.EntityFrameworkCore.DbLoggerCategory;
  31. namespace HealthMonitor.Service.Biz.db
  32. {
  33. public class TDengineService
  34. {
  35. private readonly ILogger<TDengineService> _logger;
  36. private readonly HttpHelper _httpHelper=default!;
  37. private readonly TDengineServiceConfig _configTDengineService;
  38. private readonly SqlSugarClient _clientSqlSugar;
  39. private readonly FhrPhrMapCacheManager _mgrFhrPhrMapCache;
  40. private readonly DeviceCacheManager _deviceCacheMgr;
  41. public TDengineService(ILogger<TDengineService> logger,
  42. IOptions<TDengineServiceConfig> configTDengineService,
  43. HttpHelper httpHelper,
  44. FhrPhrMapCacheManager fhrPhrMapCacheManager, DeviceCacheManager deviceCacheMgr
  45. )
  46. {
  47. _logger = logger;
  48. _configTDengineService = configTDengineService.Value;
  49. _httpHelper = httpHelper;
  50. _mgrFhrPhrMapCache = fhrPhrMapCacheManager;
  51. _deviceCacheMgr = deviceCacheMgr;
  52. _clientSqlSugar = new SqlSugarClient(new ConnectionConfig()
  53. {
  54. DbType = SqlSugar.DbType.TDengine,
  55. ConnectionString = $"Host={_configTDengineService.Host};Port={_configTDengineService.Port};Username={_configTDengineService.UserName};Password={_configTDengineService.Password};Database={_configTDengineService.DB};TsType=config_ns",
  56. IsAutoCloseConnection = true,
  57. AopEvents = new AopEvents
  58. {
  59. OnLogExecuting = (sql, p) =>
  60. {
  61. Console.WriteLine(SqlSugar.UtilMethods.GetNativeSql(sql, p));
  62. }
  63. },
  64. ConfigureExternalServices = new ConfigureExternalServices()
  65. {
  66. EntityService = (property, column) =>
  67. {
  68. if (column.SqlParameterDbType == null)
  69. {
  70. column.SqlParameterDbType = typeof(CommonPropertyConvert);
  71. }
  72. }
  73. },
  74. MoreSettings=new ConnMoreSettings()
  75. {
  76. PgSqlIsAutoToLower = false,
  77. PgSqlIsAutoToLowerCodeFirst = false,
  78. }
  79. });
  80. }
  81. public IntPtr Connection()
  82. {
  83. string host = _configTDengineService.Host;
  84. string user = _configTDengineService.UserName;
  85. string db = _configTDengineService.DB;
  86. short port = _configTDengineService.Port;
  87. string password = _configTDengineService.Password;
  88. //#if DEBUG
  89. // //string configDir = "C:/TDengine/cfg";
  90. // //TDengine.Options((int)TDengineInitOption.TSDB_OPTION_CONFIGDIR, configDir);
  91. // TDengine.Options((int)TDengineInitOption.TSDB_OPTION_TIMEZONE, "Asia/Beijing");
  92. //#endif
  93. IntPtr conn = TDengine.Connect(host, user, password, db, port);
  94. if (conn == IntPtr.Zero)
  95. {
  96. _logger.LogError($"连接 TDengine 失败....");
  97. }
  98. else
  99. {
  100. _logger.LogInformation($"连接 TDengine 成功....");
  101. }
  102. return conn;
  103. }
  104. public void ExecuteSQL(IntPtr conn, string sql)
  105. {
  106. IntPtr res = TDengine.Query(conn, sql);
  107. // Check if query success
  108. if ((res == IntPtr.Zero) || (TDengine.ErrorNo(res) != 0))
  109. {
  110. Console.Write(sql + " failure, ");
  111. // Get error message while Res is a not null pointer.
  112. if (res != IntPtr.Zero)
  113. {
  114. Console.Write("reason:" + TDengine.Error(res));
  115. }
  116. }
  117. else
  118. {
  119. Console.Write(sql + " success, {0} rows affected", TDengine.AffectRows(res));
  120. //... do something with res ...
  121. // Important: need to free result to avoid memory leak.
  122. TDengine.FreeResult(res);
  123. }
  124. }
  125. public void ExecuteQuerySQL(IntPtr conn, string sql)
  126. {
  127. IntPtr res = TDengine.Query(conn, sql);
  128. // Check if query success
  129. if ((res == IntPtr.Zero) || (TDengine.ErrorNo(res) != 0))
  130. {
  131. Console.Write(sql + " failure, ");
  132. // Get error message while Res is a not null pointer.
  133. if (res != IntPtr.Zero)
  134. {
  135. Console.Write("reason:" + TDengine.Error(res));
  136. }
  137. }
  138. else
  139. {
  140. Console.Write(sql + " success, {0} rows affected", TDengine.AffectRows(res));
  141. //... do something with res ...
  142. List<TDengineDriver.TDengineMeta> resMeta = LibTaos.GetMeta(res);
  143. List<object> resData = LibTaos.GetData(res);
  144. foreach (var meta in resMeta)
  145. {
  146. _logger.LogInformation("\t|{meta.name} {meta.TypeName()} ({meta.size})\t|", meta.name, meta.TypeName(), meta.size);
  147. }
  148. for (int i = 0; i < resData.Count; i++)
  149. {
  150. _logger.LogInformation($"|{resData[i].ToString()} \t");
  151. if (((i + 1) % resMeta.Count == 0))
  152. {
  153. _logger.LogInformation("");
  154. }
  155. }
  156. // Important: need to free result to avoid memory leak.
  157. TDengine.FreeResult(res);
  158. }
  159. }
  160. public void CheckRes(IntPtr conn, IntPtr res, String errorMsg)
  161. {
  162. if (TDengine.ErrorNo(res) != 0)
  163. {
  164. throw new Exception($"{errorMsg} since: {TDengine.Error(res)}");
  165. }
  166. }
  167. public void ExecuteInsertSQL(string sql)
  168. {
  169. var conn = Connection();
  170. try
  171. {
  172. //sql = "INSERT INTO d1001 USING meters TAGS('California.SanFrancisco', 2) VALUES ('2018-10-03 14:38:05.000', 10.30000, 219, 0.31000) ('2018-10-03 14:38:15.000', 12.60000, 218, 0.33000) ('2018-10-03 14:38:16.800', 12.30000, 221, 0.31000) " +
  173. // "d1002 USING power.meters TAGS('California.SanFrancisco', 3) VALUES('2018-10-03 14:38:16.650', 10.30000, 218, 0.25000) " +
  174. // "d1003 USING power.meters TAGS('California.LosAngeles', 2) VALUES('2018-10-03 14:38:05.500', 11.80000, 221, 0.28000)('2018-10-03 14:38:16.600', 13.40000, 223, 0.29000) " +
  175. // "d1004 USING power.meters TAGS('California.LosAngeles', 3) VALUES('2018-10-03 14:38:05.000', 10.80000, 223, 0.29000)('2018-10-03 14:38:06.500', 11.50000, 221, 0.35000)";
  176. //#if DEBUG
  177. // //string configDir = "C:/TDengine/cfg";
  178. // //TDengine.Options((int)TDengineInitOption.TSDB_OPTION_CONFIGDIR, configDir);
  179. // TDengine.Options((int)TDengineInitOption.TSDB_OPTION_TIMEZONE, "Asia/Beijing");
  180. //#endif
  181. _logger.LogInformation($"Insert SQL: {sql}");
  182. IntPtr res = TDengine.Query(conn, sql);
  183. CheckRes(conn, res, "failed to insert data");
  184. int affectedRows = TDengine.AffectRows(res);
  185. _logger.LogInformation("affectedRows {affectedRows}" , affectedRows);
  186. TDengine.FreeResult(res);
  187. }
  188. finally
  189. {
  190. TDengine.Close(conn);
  191. }
  192. }
  193. #region TDengine.Connector async query
  194. public void QueryCallback(IntPtr param, IntPtr taosRes, int code)
  195. {
  196. if (code == 0 && taosRes != IntPtr.Zero)
  197. {
  198. FetchRawBlockAsyncCallback fetchRowAsyncCallback = new FetchRawBlockAsyncCallback(FetchRawBlockCallback);
  199. TDengine.FetchRawBlockAsync(taosRes, fetchRowAsyncCallback, param);
  200. }
  201. else
  202. {
  203. _logger.LogInformation("async query data failed, failed code {code}",code);
  204. }
  205. }
  206. // Iteratively call this interface until "numOfRows" is no greater than 0.
  207. public void FetchRawBlockCallback(IntPtr param, IntPtr taosRes, int numOfRows)
  208. {
  209. if (numOfRows > 0)
  210. {
  211. _logger.LogInformation("{numOfRows} rows async retrieved", numOfRows);
  212. IntPtr pdata = TDengine.GetRawBlock(taosRes);
  213. List<TDengineMeta> metaList = TDengine.FetchFields(taosRes);
  214. List<object> dataList = LibTaos.ReadRawBlock(pdata, metaList, numOfRows);
  215. for (int i = 0; i < metaList.Count; i++)
  216. {
  217. _logger.LogInformation("{0} {1}({2}) \t|", metaList[i].name, metaList[i].type, metaList[i].size);
  218. }
  219. _logger.LogInformation("");
  220. for (int i = 0; i < dataList.Count; i++)
  221. {
  222. if (i != 0 && i % metaList.Count == 0)
  223. {
  224. _logger.LogInformation("{dataList[i]}\t|", dataList[i]);
  225. }
  226. _logger.LogInformation("{dataList[i]}\t|", dataList[i]);
  227. }
  228. TDengine.FetchRawBlockAsync(taosRes, FetchRawBlockCallback, param);
  229. }
  230. else
  231. {
  232. if (numOfRows == 0)
  233. {
  234. _logger.LogInformation("async retrieve complete.");
  235. }
  236. else
  237. {
  238. _logger.LogInformation("FetchRawBlockCallback callback error, error code {numOfRows}", numOfRows);
  239. }
  240. TDengine.FreeResult(taosRes);
  241. }
  242. }
  243. public void ExecuteQueryAsync(string sql)
  244. {
  245. var conn = Connection();
  246. QueryAsyncCallback queryAsyncCallback = new QueryAsyncCallback(QueryCallback);
  247. TDengine.QueryAsync(conn, sql, queryAsyncCallback, IntPtr.Zero);
  248. }
  249. //public void ExecuteQuery(string sql)
  250. //{
  251. // var conn = Connection();
  252. // QueryAsyncCallback queryAsyncCallback = new QueryAsyncCallback(QueryCallback);
  253. // TDengine.QueryAsync(conn, sql, queryAsyncCallback, IntPtr.Zero);
  254. //}
  255. public Aggregate GetAggregateValue(string field, string tbName, string? condition)
  256. {
  257. List<int> data = new();
  258. var sql = $"SELECT MAX({field}), MIN({field}) FROM {_configTDengineService.DB}.{tbName} WHERE {condition}";
  259. var conn = Connection();
  260. try
  261. {
  262. IntPtr res = TDengine.Query(conn, sql);
  263. // Check if query success
  264. if ((res == IntPtr.Zero) || (TDengine.ErrorNo(res) != 0))
  265. {
  266. Console.Write(sql + " failure, ");
  267. // Get error message while Res is a not null pointer.
  268. if (res != IntPtr.Zero)
  269. {
  270. Console.Write("reason:" + TDengine.Error(res));
  271. }
  272. }
  273. else
  274. {
  275. Console.Write(sql + " success, {0} rows affected", TDengine.AffectRows(res));
  276. //... do something with res ...
  277. List<TDengineMeta> resMeta = LibTaos.GetMeta(res);
  278. List<object> resData = LibTaos.GetData(res);
  279. foreach (var meta in resMeta)
  280. {
  281. Console.Write($"\t|{meta.name} {meta.TypeName()} ({meta.size})\t|");
  282. }
  283. resData.ForEach(x => data.Add(SafeType.SafeInt(x)));
  284. // Important: need to free result to avoid memory leak.
  285. TDengine.FreeResult(res);
  286. }
  287. }
  288. finally
  289. {
  290. TDengine.Close(conn);
  291. }
  292. return new Aggregate
  293. {
  294. Max = data.Count.Equals(0) ? 0 : data[0],
  295. Min = data.Count.Equals(0) ? 0 : data[1],
  296. };
  297. }
  298. public int GetAvgExceptMaxMinValue(string field, string tbName, string? condition)
  299. {
  300. List<int> data = new();
  301. var sql = $"SELECT MAX({field}), MIN({field}) FROM {_configTDengineService.DB}.{tbName} WHERE {condition}";
  302. var aggregate= GetAggregateValue(field, tbName, condition);
  303. var sqlAvg = $"SELECT AVG({field}) FROM {_configTDengineService.DB}.{tbName} WHERE {condition} AND {field} < {aggregate.Max} and {field} > {aggregate.Min}";
  304. var conn = Connection();
  305. try
  306. {
  307. IntPtr res = TDengine.Query(conn, sqlAvg);
  308. // Check if query success
  309. if ((res == IntPtr.Zero) || (TDengine.ErrorNo(res) != 0))
  310. {
  311. Console.Write(sqlAvg + " failure, ");
  312. // Get error message while Res is a not null pointer.
  313. if (res != IntPtr.Zero)
  314. {
  315. Console.Write("reason:" + TDengine.Error(res));
  316. }
  317. }
  318. else
  319. {
  320. Console.Write(sqlAvg + " success, {0} rows affected", TDengine.AffectRows(res));
  321. //... do something with res ...
  322. List<TDengineMeta> resMeta = LibTaos.GetMeta(res);
  323. List<object> resData = LibTaos.GetData(res);
  324. foreach (var meta in resMeta)
  325. {
  326. Console.Write($"\t|{meta.name} {meta.TypeName()} ({meta.size})\t|");
  327. }
  328. resData.ForEach(x => data.Add(SafeType.SafeInt(x)));
  329. // Important: need to free result to avoid memory leak.
  330. TDengine.FreeResult(res);
  331. }
  332. }
  333. finally
  334. {
  335. TDengine.Close(conn);
  336. }
  337. return data.Count.Equals(0) ? 0 : data[0];
  338. }
  339. #endregion
  340. #region RestAPI
  341. public async Task<string?> ExecuteQuerySQLRestResponse(string sql)
  342. {
  343. var url = $"http://{_configTDengineService.Host}:{_configTDengineService.RestPort}/rest/sql/{_configTDengineService.DB}";
  344. List<KeyValuePair<string, string>> headers = new()
  345. {
  346. new KeyValuePair<string, string>("Authorization", "Basic " + _configTDengineService.Token)
  347. };
  348. var result = await _httpHelper.HttpToPostAsync(url, sql, headers).ConfigureAwait(false);
  349. return result;
  350. }
  351. public async Task<string?> ExecuteSelectRestResponseAsync( string tbName, string condition="1", string field = "*")
  352. {
  353. var url = $"http://{_configTDengineService.Host}:{_configTDengineService.RestPort}/rest/sql/{_configTDengineService.DB}";
  354. var sql = $"SELECT {field} FROM {_configTDengineService.DB}.{tbName} WHERE {condition}";
  355. List<KeyValuePair<string, string>> headers = new()
  356. {
  357. new KeyValuePair<string, string>("Authorization", "Basic " + _configTDengineService.Token)
  358. };
  359. _logger.LogInformation($"{nameof(ExecuteSelectRestResponseAsync)} --- SQL 语句执行 {sql}");
  360. var result = await _httpHelper.HttpToPostAsync(url, sql, headers).ConfigureAwait(false);
  361. return result;
  362. }
  363. public async Task<bool> GernalRestSql(string sql)
  364. {
  365. //"http://{server}:{port}/rest/sql/{db}"
  366. var url = $"http://{_configTDengineService.Host}:{_configTDengineService.RestPort}/rest/sql/{_configTDengineService.DB}";
  367. List<KeyValuePair<string, string>> headers = new()
  368. {
  369. new KeyValuePair<string, string>("Authorization", "Basic " + _configTDengineService.Token)
  370. };
  371. var result = await _httpHelper.HttpToPostAsync(url, sql, headers).ConfigureAwait(false);
  372. var res = JsonConvert.DeserializeObject<TDengineRestResBase>(result!);
  373. if (result != null)
  374. {
  375. if (res?.Code == 0)
  376. {
  377. _logger.LogInformation($"{nameof(GernalRestSql)},SQL 语句执行成功|{sql}");
  378. return true;
  379. }
  380. else
  381. {
  382. _logger.LogWarning($"{nameof(GernalRestSql)},SQL 语句执行失败||{sql}");
  383. return false;
  384. }
  385. }
  386. else
  387. {
  388. _logger.LogError($"{nameof(GernalRestSql)},TDengine 服务器IP:{_configTDengineService.Host} 错误,请联系运维人员");
  389. return false;
  390. }
  391. //return res.Code==0;
  392. }
  393. public async Task<string?> GernalRestSqlResTextAsync(string sql)
  394. {
  395. _logger.LogInformation($"执行 SQL: {nameof(GernalRestSqlResTextAsync)}--{sql}");
  396. var url = $"http://{_configTDengineService.Host}:{_configTDengineService.RestPort}/rest/sql/{_configTDengineService.DB}";
  397. List<KeyValuePair<string, string>> headers = new()
  398. {
  399. new KeyValuePair<string, string>("Authorization", "Basic " + _configTDengineService.Token)
  400. };
  401. var result = await _httpHelper.HttpToPostAsync(url, sql, headers).ConfigureAwait(false);
  402. return result;
  403. }
  404. /** 取消
  405. /// <summary>
  406. /// 最大值,最小值聚合
  407. /// </summary>
  408. /// <param name="field"></param>
  409. /// <param name="tbName"></param>
  410. /// <param name="condition"></param>
  411. /// <returns></returns>
  412. public async Task<Aggregate> GetAggregateValueAsync(string field,string tbName,string? condition)
  413. {
  414. var sql = $"SELECT MAX({field}), MIN({field}) FROM {_configTDengineService.DB}.{tbName} WHERE {condition}";
  415. var result = await GernalRestSqlResTextAsync(sql);
  416. var res = JsonConvert.DeserializeObject<Aggregate>(result!);
  417. List<dynamic> data = res?.Data!;
  418. return new Aggregate
  419. {
  420. Max = data.Count.Equals(0) ? 0 : data[0][0],
  421. Min = data.Count.Equals(0) ? 0 : data[0][1],
  422. };
  423. }
  424. /// <summary>
  425. /// 去除最大值和最小值后的平均值
  426. /// </summary>
  427. /// <param name="field"></param>
  428. /// <param name="tbName"></param>
  429. /// <param name="condition"></param>
  430. /// <returns></returns>
  431. public async Task<int> GetAvgExceptMaxMinValueAsync(string field, string tbName, string? condition)
  432. {
  433. var sql = $"SELECT MAX({field}), MIN({field}) FROM {_configTDengineService.DB}.{tbName} WHERE {condition}";
  434. var result = await GernalRestSqlResTextAsync(sql);
  435. var res = JsonConvert.DeserializeObject<TDengineRestResBase>(result!);
  436. _logger.LogInformation($"最大小值:{sql}");
  437. List<dynamic> data = res?.Data!;
  438. var sqlAvg = $"SELECT AVG({field}) FROM {_configTDengineService.DB}.{tbName} WHERE {condition} AND {field} < { (data.Count.Equals(0)? 0: data[0][0]) } and {field} > {(data.Count.Equals(0) ? 0 : data[0][1])}";
  439. result = await GernalRestSqlResTextAsync(sqlAvg);
  440. _logger.LogInformation($"sqlAvg:{sqlAvg}");
  441. res = JsonConvert.DeserializeObject<TDengineRestResBase>(result!);
  442. data = res?.Data!;
  443. return data.Count.Equals(0)?0:(int)data[0][0];
  444. }
  445. /// <summary>
  446. /// 获取最后的记录
  447. /// </summary>
  448. /// <param name="tbName"></param>
  449. /// <param name="condition"></param>
  450. /// <returns></returns>
  451. public async Task<JArray?> GetLastAsync(string tbName, string? condition)
  452. {
  453. var sql = $"SELECT last_row(*) FROM {_configTDengineService.DB}.{tbName} WHERE {condition}";
  454. var result = await GernalRestSqlResTextAsync(sql);
  455. var res = JsonConvert.DeserializeObject<TDengineRestResBase>(result!);
  456. if(res?.Data?.Count==0) return null;
  457. List<dynamic> data = res?.Data!;
  458. return data[0] as JArray;
  459. }
  460. public async Task<int> GetCount(string tbName, string? condition)
  461. {
  462. var sql = $"SELECT count(ts) FROM {_configTDengineService.DB}.{tbName} WHERE {condition}";
  463. var result = await GernalRestSqlResTextAsync(sql);
  464. var res = JsonConvert.DeserializeObject<TDengineRestResBase>(result!);
  465. List<dynamic> data = res?.Data!;
  466. return SafeType.SafeInt(data[0][0]);
  467. }
  468. */
  469. #endregion
  470. /// <summary>
  471. /// 平均值算法(去除最大值,最小值和大于标定值的平均值)
  472. /// </summary>
  473. /// <param name="numToRemove"></param>
  474. /// <param name="collection"></param>
  475. /// <param name="max"></param>
  476. /// <param name="min"></param>
  477. /// <returns></returns>
  478. public static decimal AverageAfterRemovingOneMinMaxRef(List<int> collection, int max, int min,int refValue)
  479. {
  480. collection.Remove(max);
  481. collection.Remove(min);
  482. collection.RemoveAll(_ => _ > refValue);
  483. if (collection.Count < 2)
  484. {
  485. throw new ArgumentException($"数据集{collection.ToArray()},去掉一个最大值 {max}和一个最小值{min},异常值(大于标定值{refValue}),后数据值不足");
  486. }
  487. return (decimal)collection.Average(x => x);
  488. //var values = ParseData.Select(valueSelector).ToList();
  489. //collection = values.Select(i => (int)i).ToArray();
  490. //if (values.Count <= 2)
  491. //{
  492. // throw new ArgumentException("Not enough elements to remove.");
  493. //}
  494. //// Remove the specified number of minimum and maximum values
  495. ////values.RemoveAll(_ => _ == values.Min());
  496. ////values.RemoveAll(_ => _ == values.Max());
  497. //max = (int)values.Max();
  498. //min = (int)values.Min();
  499. //values.Remove(max);
  500. //values.Remove(min);
  501. //// Remove values less than the specified threshold
  502. //values.RemoveAll(_ => _ > numToRemove);
  503. //// Calculate and return the average of the remaining values
  504. //return values.Average();
  505. }
  506. /// <summary>
  507. /// 去除最大值和最小值各一个(列表的头和尾),再去除异常值
  508. /// </summary>
  509. /// <param name="systolicRefValue"></param>
  510. /// <param name="hmBpParser"></param>
  511. /// <returns></returns>
  512. public decimal[] AverageAfterRemovingOneMinMaxRef(int systolicRefValue, ParseTDengineRestResponse<BloodPressureModel>? hmBpParser)
  513. {
  514. var sortedList = hmBpParser?.Select(i => i)
  515. .Where(i => i.IsDisplay.Equals(true))
  516. .OrderByDescending(i => i.SystolicValue)
  517. .ThenByDescending(i => i.DiastolicValue)
  518. .ToList();
  519. //_logger.LogInformation($"计算时间段排列数据集:{JsonConvert.SerializeObject(sortedList)}");
  520. // 去除最大值和最小值各一个(列表的头和尾)
  521. var trimmedList = sortedList?
  522. .Skip(1)
  523. .Take(sortedList.Count - 2)
  524. .ToList();
  525. //_logger.LogInformation($"计算去除最大值和最小值各一个数据集:{JsonConvert.SerializeObject(trimmedList)}");
  526. // 去除异常值
  527. var filteredList = trimmedList?.Where(bp => bp.SystolicValue < SafeType.SafeInt(systolicRefValue!)).ToList();
  528. //_logger.LogInformation($"计算除异常值个数据集:{JsonConvert.SerializeObject(filteredList)}");
  529. if (filteredList?.Count < 2)
  530. {
  531. // throw new ArgumentException("数据不够不能计算");
  532. // 平均值为0,说明数据不足,不能计算增量值
  533. return new decimal[] { 0M, 0M };
  534. }
  535. var systolicAvg = filteredList?.Select(bp => bp.SystolicValue).Average();
  536. var diastolicAvg = filteredList?.Select(bp => bp.DiastolicValue).Average();
  537. return new decimal[] { (decimal)systolicAvg!, (decimal)diastolicAvg! };
  538. }
  539. /// <summary>
  540. /// 去除最大值和最小值各一个(列表的头和尾)
  541. /// </summary>
  542. /// <param name="systolicRefValue"></param>
  543. /// <param name="hmBpParser"></param>
  544. /// <returns></returns>
  545. public decimal[] AverageAfterRemovingOneMinMaxRef(ParseTDengineRestResponse<BloodPressureModel>? hmBpParser)
  546. {
  547. var sortedList = hmBpParser?.Select(i => i)
  548. .Where(i => i.IsDisplay.Equals(true))
  549. .OrderByDescending(i => i.SystolicValue)
  550. .ThenByDescending(i => i.DiastolicValue)
  551. .ToList();
  552. //_logger.LogInformation($"计算时间段排列数据集:{JsonConvert.SerializeObject(sortedList)}");
  553. // 去除最大值和最小值各一个(列表的头和尾)
  554. var trimmedList = sortedList?
  555. .Skip(1)
  556. .Take(sortedList.Count - 2)
  557. .ToList();
  558. //_logger.LogInformation($"计算去除最大值和最小值各一个数据集:{JsonConvert.SerializeObject(trimmedList)}");
  559. var filteredList = trimmedList?.ToList();
  560. //_logger.LogInformation($"计算除异常值个数据集:{JsonConvert.SerializeObject(filteredList)}");
  561. if (filteredList?.Count < 2)
  562. {
  563. // throw new ArgumentException("数据不够不能计算");
  564. // 平均值为0,说明数据不足,不能计算增量值
  565. return new decimal[] { 0M, 0M };
  566. }
  567. var systolicAvg = filteredList?.Select(bp => bp.SystolicValue).Average();
  568. var diastolicAvg = filteredList?.Select(bp => bp.DiastolicValue).Average();
  569. return new decimal[] { (decimal)systolicAvg!, (decimal)diastolicAvg! };
  570. }
  571. #region SqlSugarClient
  572. /**
  573. public async Task InsertFetalHeartRateAsync()
  574. {
  575. var tableName = typeof(FetalHeartRateModel)
  576. .GetCustomAttribute<STableAttribute>()?
  577. .STableName;
  578. if (tableName == null)
  579. {
  580. throw new InvalidOperationException("STableAttribute not found on FetalHeartRateModel class.");
  581. }
  582. await _clientSqlSugar.Ado.ExecuteCommandAsync($"create table IF NOT EXISTS hm_fhr_00 using {tableName} tags('00')");
  583. _clientSqlSugar.Insertable(new FetalHeartRateModel()
  584. {
  585. Timestamp = DateTime.Now,
  586. CreateTime = DateTime.Now,
  587. FetalHeartRate = 90,
  588. FetalHeartRateId = Guid.NewGuid().ToString("D"),
  589. IsDisplay = false,
  590. Method = 1,
  591. PersonId = Guid.NewGuid().ToString("D"),
  592. MessageId = Guid.NewGuid().ToString("D"),
  593. SerialNumber = Guid.NewGuid().ToString("D"),
  594. DeviceKey = Guid.NewGuid().ToString("D"),
  595. SerialTailNumber = "00",
  596. LastUpdate = DateTime.Now,
  597. }).AS("hm_fhr_00").ExecuteCommand();
  598. }
  599. public Task<FetalHeartRateModel> GetFirst()
  600. {
  601. var tableName = typeof(FetalHeartRateModel)
  602. .GetCustomAttribute<STableAttribute>()?
  603. .STableName;
  604. var first = _clientSqlSugar
  605. .Queryable<FetalHeartRateModel>()
  606. .AS(tableName)
  607. .OrderByDescending(x => x.Timestamp).FirstAsync();
  608. return first;
  609. }
  610. */
  611. /// <summary>
  612. /// 插入记录
  613. /// </summary>
  614. /// <typeparam name="T"></typeparam>
  615. /// <param name="model"></param>
  616. /// <param name="tbName">子表名称</param>
  617. /// <returns></returns>
  618. /// <exception cref="InvalidOperationException"></exception>
  619. public async Task InsertAsync<T>(string tbName, T model)
  620. {
  621. var stbName = typeof(T)
  622. .GetCustomAttribute<STableAttribute>()?
  623. .STableName;
  624. if (stbName == null)
  625. {
  626. throw new InvalidOperationException($"STableAttribute not found on {nameof(T)} class.");
  627. }
  628. var tailNo = typeof(T).GetProperty("SerialTailNumber")?.GetValue(model)?.ToString();
  629. if (string.IsNullOrEmpty(tailNo))
  630. {
  631. throw new InvalidOperationException($"SerialNumberAttribute not found on {nameof(T)} class.");
  632. }
  633. var tbFullName = $"{tbName}_{tailNo}";
  634. await _clientSqlSugar.Ado.ExecuteCommandAsync($"create table IF NOT EXISTS {tbFullName} using {stbName} tags('{tailNo}')");
  635. //_clientSqlSugar.InsertableByDynamic(model).AS(tbFullName).ExecuteCommand();
  636. _clientSqlSugar.InsertableByObject(model).AS(tbFullName).ExecuteCommand();
  637. }
  638. public Task<T> GetLastAsync<T>() where T : class
  639. {
  640. var tableName = typeof(T)
  641. .GetCustomAttribute<STableAttribute>()?
  642. .STableName;
  643. // 创建一个表示 Timestamp 属性的表达式
  644. var parameter = Expression.Parameter(typeof(T), "x");
  645. var property = Expression.Property(parameter, "Timestamp");
  646. var lambda = Expression.Lambda<Func<T, object>>(Expression.Convert(property, typeof(object)), parameter);
  647. var first = _clientSqlSugar
  648. .Queryable<T>()
  649. .AS(tableName)
  650. .OrderByDescending(lambda).FirstAsync();
  651. return first;
  652. }
  653. public Task<T> GetLastAsync<T>(string serialNo) where T : class
  654. {
  655. var tableName = typeof(T)
  656. .GetCustomAttribute<STableAttribute>()?
  657. .STableName;
  658. // 创建表示 Timestamp 属性的表达式
  659. var parameter = Expression.Parameter(typeof(T), "x");
  660. var timestampProperty = Expression.Property(parameter, "Timestamp");
  661. var timestampLambda = Expression.Lambda<Func<T, object>>(Expression.Convert(timestampProperty, typeof(object)), parameter);
  662. // 创建表示 SerialNo 属性的表达式
  663. var serialNoProperty = Expression.Property(parameter, "SerialNumber");
  664. var serialNoConstant = Expression.Constant(serialNo);
  665. var equalExpression = Expression.Equal(serialNoProperty, serialNoConstant);
  666. var serialNoLambda = Expression.Lambda<Func<T, bool>>(equalExpression, parameter);
  667. var first = _clientSqlSugar
  668. .Queryable<T>()
  669. .AS(tableName)
  670. .Where(serialNoLambda)
  671. .OrderByDescending(timestampLambda).FirstAsync();
  672. return first;
  673. }
  674. public async Task<List<T>> GetBySerialNoAsync<T>(string serialNo) where T : class
  675. {
  676. var tableName = typeof(T)
  677. .GetCustomAttribute<STableAttribute>()?
  678. .STableName;
  679. // 创建表示 Timestamp 属性的表达式
  680. var parameter = Expression.Parameter(typeof(T), "x");
  681. var timestampProperty = Expression.Property(parameter, "Timestamp");
  682. var timestampLambda = Expression.Lambda<Func<T, object>>(Expression.Convert(timestampProperty, typeof(object)), parameter);
  683. // 创建表示 SerialNo 属性的表达式
  684. var serialNoProperty = Expression.Property(parameter, "SerialNumber");
  685. var serialNoConstant = Expression.Constant(serialNo);
  686. var equalExpression = Expression.Equal(serialNoProperty, serialNoConstant);
  687. var serialNoLambda = Expression.Lambda<Func<T, bool>>(equalExpression, parameter);
  688. var res = await _clientSqlSugar
  689. .Queryable<T>()
  690. .AS(tableName)
  691. .Where(serialNoLambda)
  692. .OrderByDescending(timestampLambda).ToListAsync();
  693. return res;
  694. }
  695. public async Task<List<T>> GetBySerialNoAsync<T>(string serialNo, int days) where T : class
  696. {
  697. var tableName = typeof(T)
  698. .GetCustomAttribute<STableAttribute>()?
  699. .STableName;
  700. // 创建表示 Timestamp 属性的表达式
  701. var parameter = Expression.Parameter(typeof(T), "x");
  702. var timestampProperty = Expression.Property(parameter, "Timestamp");
  703. var timestampLambda = Expression.Lambda<Func<T, object>>(Expression.Convert(timestampProperty, typeof(object)), parameter);
  704. // 创建表示 SerialNo 属性的表达式
  705. var serialNoProperty = Expression.Property(parameter, "SerialNumber");
  706. var serialNoConstant = Expression.Constant(serialNo);
  707. var equalExpression = Expression.Equal(serialNoProperty, serialNoConstant);
  708. // 创建表示 Timestamp 大于指定天数的表达式
  709. var daysAgo = DateTime.Now.AddDays(-days);
  710. var daysAgoConstant = Expression.Constant(daysAgo, typeof(DateTime));
  711. var greaterThanExpression = Expression.GreaterThan(timestampProperty, daysAgoConstant);
  712. // 合并 SerialNo 和 Timestamp 条件
  713. var combinedExpression = Expression.AndAlso(equalExpression, greaterThanExpression);
  714. var combinedLambda = Expression.Lambda<Func<T, bool>>(combinedExpression, parameter);
  715. var res = await _clientSqlSugar
  716. .Queryable<T>()
  717. .AS(tableName)
  718. .Where(combinedLambda)
  719. .OrderByDescending(timestampLambda).ToListAsync();
  720. return res;
  721. }
  722. public async Task DeleteAllBySerialNoCMDAsync<T>(string serialNo) where T : class
  723. {
  724. var records = await GetBySerialNoAsync<T>(serialNo, 365);
  725. var stbName = typeof(T)
  726. .GetCustomAttribute<STableAttribute>()?
  727. .STableName;
  728. var tasks = records.Select(async r =>
  729. {
  730. Type modelType = typeof(T);
  731. PropertyInfo timestampProperty = typeof(T).GetProperty("Timestamp")!;
  732. object timestampValue = timestampProperty.GetValue(r)!;
  733. var ts = ((DateTime)timestampValue);
  734. var startTimestamp = ts.ToString("yyyy-MM-dd HH:mm:ss.fff");
  735. var endTimestamp = ts.AddMilliseconds(1).ToString("yyyy-MM-dd HH:mm:ss.fff");
  736. var sql = $"DELETE FROM {stbName} WHERE ts >= '{startTimestamp}' AND ts < '{endTimestamp}'";
  737. var res= await _clientSqlSugar.Ado.ExecuteCommandAsync(sql);
  738. Console.WriteLine(res);
  739. });
  740. await Task.WhenAll(tasks);
  741. }
  742. #endregion
  743. #region 胎心算法
  744. ///// <summary>
  745. ///// 计算个人一般心率(最大值,最大值,最小值)
  746. ///// </summary>
  747. ///// <param name="serialNo"></param>
  748. ///// <param name="days"></param>
  749. ///// <param name="percentage"></param>
  750. ///// <returns></returns>
  751. //public async Task<PregnancyCommonHeartRateModel?> InitPregnancyCommonHeartRateModeAsync(string serialNo, int days = 7, int percentage = 90)
  752. //{
  753. // var tableName = typeof(PregnancyHeartRateModel)
  754. // .GetCustomAttribute<STableAttribute>()?
  755. // .STableName;
  756. // var daysAgo = DateTime.Now.AddDays(-days);
  757. // var collection = await _clientSqlSugar
  758. // .Queryable<PregnancyHeartRateModel>()
  759. // .AS(tableName)
  760. // .Where(i => i.SerialNumber.Equals(serialNo))
  761. // .Where(i => i.Timestamp > daysAgo)
  762. // .OrderByDescending(i => i.Timestamp)
  763. // .ToArrayAsync();
  764. // var res = collection
  765. // .Select(i => i.PregnancyHeartRate).ToList();
  766. // // 心率数据量必须30个以上才进行计算
  767. // if (res.Count < 30)
  768. // {
  769. // _logger.LogInformation($"{serialNo} 心率数据不足,无法计算其众数");
  770. // return null;
  771. // }
  772. // #region 计算众数
  773. // var mode = res.GroupBy(n => n)
  774. // .OrderByDescending(g => g.Count())
  775. // .First()
  776. // .Key;
  777. // Console.WriteLine("众数是: " + mode);
  778. // // 如果有多个众数的情况
  779. // var maxCount = res.GroupBy(n => n)
  780. // .Max(g => g.Count());
  781. // var modes = res.GroupBy(n => n)
  782. // .Where(g => g.Count() == maxCount)
  783. // .Select(g => g.Key)
  784. // .ToList();
  785. // // 多个众数,选择最接近平均数或中位数的众数
  786. // if (modes.Count > 1)
  787. // {
  788. // // 计算平均值
  789. // double average = res.Average();
  790. // Console.WriteLine("平均值是: " + average);
  791. // // 计算中位数
  792. // double median;
  793. // int count = res.Count;
  794. // var sortedRes = res.OrderBy(n => n).ToList();
  795. // if (count % 2 == 0)
  796. // {
  797. // // 偶数个元素,取中间两个数的平均值
  798. // median = (sortedRes[count / 2 - 1] + sortedRes[count / 2]) / 2.0;
  799. // }
  800. // else
  801. // {
  802. // // 奇数个元素,取中间的数
  803. // median = sortedRes[count / 2];
  804. // }
  805. // //Console.WriteLine("中位数是: " + median);
  806. // _logger.LogInformation($"{serialNo} 中位数是: " + median);
  807. // // 找出最接近平均值的众数
  808. // //var closestToAverage = modes.OrderBy(m => Math.Abs(m - average)).First();
  809. // //Console.WriteLine("最接近平均值的众数是: " + closestToAverage);
  810. // // 找出最接近中位数的众数
  811. // var closestToMedian = modes.OrderBy(m => Math.Abs(m - median)).First();
  812. // _logger.LogInformation($"{serialNo} 最接近中位数的众数是: " + closestToMedian);
  813. // mode = closestToMedian;
  814. // }
  815. // #endregion
  816. // // 计算需要的数量
  817. // int requiredCount = (int)(res.Count * 0.9);
  818. // // 从原始数据集中获取最接近众数的元素
  819. // var closestToModeData = res.OrderBy(n => Math.Abs(n - mode))
  820. // .Take(requiredCount)
  821. // .ToList();
  822. // // 输出新数据集
  823. // _logger.LogInformation($"{serialNo} 新数据集: " + string.Join(", ", closestToModeData));
  824. // _logger.LogInformation($"{serialNo} 新数据集的数量: " + closestToModeData.Count);
  825. // var fhrMap = _mgrFhrPhrMapCache.GetHeartRatesMap();
  826. // var watchConfig = await _deviceCacheMgr.GetGpsDeviceWatchConfigCacheObjectBySerialNoAsync(serialNo, "0067");
  827. // if (watchConfig == null)
  828. // {
  829. // return null;
  830. // }
  831. // // long.TryParse(watchConfig["EDOC"]!.ToString(), out long edoc);
  832. // // "EDOC": "1720860180652",当前时间 - (EDOC - 280) days =怀孕时间
  833. // //edoc = edoc.ToString().Length == 10 ? edoc * 1000 : edoc;
  834. // var edoc = DateTimeUtil.ToDateTime(watchConfig["EDOC"]!.ToString());
  835. // int pregnancyWeek = (DateTime.Now - edoc.AddDays(-280)).Days / 7;
  836. // _logger.LogInformation($"IMEI {serialNo},EDOC:{edoc},NOW:{DateTime.Now.ToString("yyyy-MM-dd HH:mm:ss")},SinceNOW:{edoc.AddDays(-280).ToString("yyyy-MM-dd HH:mm:ss")},怀孕周数 {pregnancyWeek}");
  837. // float statMaxValueFprCoefficient = 0f;
  838. // float statMinValueFprCoefficient = 0f;
  839. // float StatModeAvgFprCoefficient = 0f;
  840. // // 20-45周之间
  841. // if (pregnancyWeek >= 12 && pregnancyWeek <= 45)
  842. // {
  843. // var map = fhrMap
  844. // .Where(i =>
  845. // i.PregnancyPeriod![0] <= pregnancyWeek &&
  846. // i.PregnancyPeriod[1] >= pregnancyWeek &&
  847. // i.PregnancyHeartRateRange![0] <= mode &&
  848. // i.PregnancyHeartRateRange[1] >= mode)
  849. // .FirstOrDefault();
  850. // if (map != null)
  851. // {
  852. // statMaxValueFprCoefficient = (float)Math.Round((decimal)map.FetalHeartRateRange![1] / res.Max(), 3);
  853. // statMinValueFprCoefficient = (float)Math.Round((decimal)map.FetalHeartRateRange[0] / res.Min(), 3);
  854. // StatModeAvgFprCoefficient = (float)Math.Round((decimal)map.FetalHeartRateAverage / mode, 3);
  855. // }
  856. // }
  857. // return new PregnancyCommonHeartRateModel()
  858. // {
  859. // Timestamp = DateTime.Now,
  860. // PersonId = collection.First().DeviceKey,
  861. // DeviceKey = collection.First().DeviceKey,
  862. // SerialNumber = collection.First().SerialNumber,
  863. // Mode = mode,
  864. // Percentage = percentage,
  865. // MaxValue = closestToModeData.Max(),
  866. // MinValue = closestToModeData.Min(),
  867. // OriginalMaxValue = res.Max(),
  868. // OriginalMinValue = res.Min(),
  869. // CreateTime = DateTime.Now,
  870. // StatStartTime = collection.OrderBy(i => i.Timestamp).Select(i => i.Timestamp).First(),
  871. // StatEndTime = collection.OrderBy(i => i.Timestamp).Select(i => i.Timestamp).Last(),
  872. // StatMaxValueFprCoefficient = statMaxValueFprCoefficient,
  873. // StatMinValueFprCoefficient = statMinValueFprCoefficient,
  874. // StatModeAvgFprCoefficient = StatModeAvgFprCoefficient,
  875. // Remark = string.Empty,
  876. // SerialTailNumber = serialNo.Substring(serialNo.Length - 2)
  877. // };
  878. //}
  879. public async Task<PregnancyCommonHeartRateModel?> InitPregnancyCommonHeartRateModeAsync(string serialNo, int days = 7, int percentage = 90, int highFreqSampleInterval=0)
  880. {
  881. var tableName = typeof(PregnancyHeartRateModel)
  882. .GetCustomAttribute<STableAttribute>()?
  883. .STableName;
  884. var daysAgo = DateTime.Now.AddDays(-days);
  885. var collection = await _clientSqlSugar
  886. .Queryable<PregnancyHeartRateModel>()
  887. .AS(tableName)
  888. .Where(i => i.SerialNumber.Equals(serialNo))
  889. .Where(i => i.LastUpdate > daysAgo)
  890. .OrderByDescending(i => i.LastUpdate)
  891. .ToListAsync();
  892. // 去除高频数据
  893. var filteredCollection = GetNonFreqPregnancyHeartRate(collection, highFreqSampleInterval);
  894. // 心率数据量必须30个以上才进行计算
  895. if (filteredCollection.Count < 30)
  896. {
  897. _logger.LogInformation($"{serialNo} 心率数据不足,无法计算其众数");
  898. return null;
  899. }
  900. var res = filteredCollection
  901. .Select(i => i.PregnancyHeartRate).ToList();
  902. //// 心率数据量必须30个以上才进行计算
  903. //if (res.Count < 30)
  904. //{
  905. // _logger.LogInformation($"{serialNo} 心率数据不足,无法计算其众数");
  906. // return null;
  907. //}
  908. var listRes = filteredCollection.Select(i => new { last_update=i.LastUpdate.ToString("yyyy-MM-dd HH:mm:ss"),heart_rate=i.PregnancyHeartRate }).ToList();
  909. //_logger.LogInformation($"highFreqSampleInterval:{highFreqSampleInterval},{serialNo} 去除高频数据后列表: {JsonConvert.SerializeObject(listRes)}");
  910. _logger.LogInformation($"{serialNo} 去除高频数据后的数据集: " + string.Join(", ", res));
  911. #region 计算众数
  912. var mode = res.GroupBy(n => n)
  913. .OrderByDescending(g => g.Count())
  914. .First()
  915. .Key;
  916. _logger.LogInformation("众数是: " + mode);
  917. // 如果有多个众数的情况
  918. var maxCount = res.GroupBy(n => n)
  919. .Max(g => g.Count());
  920. var modes = res.GroupBy(n => n)
  921. .Where(g => g.Count() == maxCount)
  922. .Select(g => g.Key)
  923. .ToList();
  924. // 多个众数,选择最接近平均数或中位数的众数
  925. if (modes.Count > 1)
  926. {
  927. // 计算平均值
  928. double average = res.Average();
  929. Console.WriteLine("平均值是: " + average);
  930. // 计算中位数
  931. double median;
  932. int count = res.Count;
  933. var sortedRes = res.OrderBy(n => n).ToList();
  934. if (count % 2 == 0)
  935. {
  936. // 偶数个元素,取中间两个数的平均值
  937. median = (sortedRes[count / 2 - 1] + sortedRes[count / 2]) / 2.0;
  938. }
  939. else
  940. {
  941. // 奇数个元素,取中间的数
  942. median = sortedRes[count / 2];
  943. }
  944. //Console.WriteLine("中位数是: " + median);
  945. _logger.LogInformation($"{serialNo} 中位数是: " + median);
  946. // 找出最接近平均值的众数
  947. //var closestToAverage = modes.OrderBy(m => Math.Abs(m - average)).First();
  948. //Console.WriteLine("最接近平均值的众数是: " + closestToAverage);
  949. // 找出最接近中位数的众数
  950. var closestToMedian = modes.OrderBy(m => Math.Abs(m - median)).First();
  951. _logger.LogInformation($"{serialNo} 最接近中位数的众数是: " + closestToMedian);
  952. mode = closestToMedian;
  953. }
  954. #endregion
  955. // 计算需要的数量
  956. int requiredCount = (int)(res.Count * 0.9);
  957. // 从原始数据集中获取最接近众数的元素
  958. var closestToModeData = res.OrderBy(n => Math.Abs(n - mode))
  959. .Take(requiredCount)
  960. .ToList();
  961. // 输出新数据集
  962. _logger.LogInformation($"{serialNo} 新数据集: " + string.Join(", ", closestToModeData));
  963. _logger.LogInformation($"{serialNo} 新数据集的数量: {closestToModeData.Count},最大值: {closestToModeData.Max()},最小值 {closestToModeData.Min()},原始最大值: {res.Max()},原始最小值 {res.Min()}" );
  964. var fhrMap = _mgrFhrPhrMapCache.GetHeartRatesMap();
  965. var watchConfig = await _deviceCacheMgr.GetGpsDeviceWatchConfigCacheObjectBySerialNoAsync(serialNo, "0067");
  966. if (watchConfig == null)
  967. {
  968. return null;
  969. }
  970. // long.TryParse(watchConfig["EDOC"]!.ToString(), out long edoc);
  971. // "EDOC": "1720860180652",当前时间 - (EDOC - 280) days =怀孕时间
  972. //edoc = edoc.ToString().Length == 10 ? edoc * 1000 : edoc;
  973. var edoc = DateTimeUtil.ToDateTime(watchConfig["EDOC"]!.ToString());
  974. int pregnancyWeek = (DateTime.Now - edoc.AddDays(-280)).Days / 7;
  975. _logger.LogInformation($"IMEI {serialNo},EDOC:{edoc},NOW:{DateTime.Now.ToString("yyyy-MM-dd HH:mm:ss")},SinceNOW:{edoc.AddDays(-280).ToString("yyyy-MM-dd HH:mm:ss")},怀孕周数 {pregnancyWeek}");
  976. float statMaxValueFprCoefficient = 0f;
  977. float statMinValueFprCoefficient = 0f;
  978. float StatModeAvgFprCoefficient = 0f;
  979. // 如果最大值与最小值在60~100范围内,都按个固定值下发高频阈值。
  980. var maxValue = closestToModeData.Max();
  981. var minValue = closestToModeData.Min();
  982. //if ((maxValue >= 60 && maxValue <= 100) && (minValue >= 60 && minValue <= 100))
  983. //{
  984. // minValue = 60;
  985. // maxValue = 100;
  986. //}
  987. // 最大值最小值两边扩散
  988. // 最小值不能大于60
  989. minValue = minValue > 60 ? 60 : minValue;
  990. // 最大值不能少于100
  991. maxValue= maxValue < 100 ? 100 : maxValue;
  992. // 20-45周之间
  993. if (pregnancyWeek >= 12 && pregnancyWeek <= 45)
  994. {
  995. var map = fhrMap
  996. .Where(i =>
  997. i.PregnancyPeriod![0] <= pregnancyWeek &&
  998. i.PregnancyPeriod[1] >= pregnancyWeek
  999. //&&i.PregnancyHeartRateRange![0] <= mode && i.PregnancyHeartRateRange[1] >= mode
  1000. )
  1001. .FirstOrDefault();
  1002. if (map != null)
  1003. {
  1004. //statMaxValueFprCoefficient = (float)Math.Round((decimal)map.FetalHeartRateRange![1] / closestToModeData.Max(), 3);
  1005. //statMinValueFprCoefficient = (float)Math.Round((decimal)map.FetalHeartRateRange[0] / closestToModeData.Min(), 3);
  1006. statMaxValueFprCoefficient = (float)Math.Round((decimal)map.FetalHeartRateRange![1] / maxValue, 3);
  1007. statMinValueFprCoefficient = (float)Math.Round((decimal)map.FetalHeartRateRange[0] / minValue, 3);
  1008. StatModeAvgFprCoefficient = (float)Math.Round((decimal)map.FetalHeartRateAverage / mode, 3);
  1009. }
  1010. }
  1011. return new PregnancyCommonHeartRateModel()
  1012. {
  1013. Timestamp = DateTime.Now,
  1014. PersonId = collection.First().DeviceKey,
  1015. DeviceKey = collection.First().DeviceKey,
  1016. SerialNumber = collection.First().SerialNumber,
  1017. Mode = mode,
  1018. Percentage = percentage,
  1019. MaxValue = maxValue,
  1020. MinValue = minValue,
  1021. OriginalMaxValue = res.Max(),
  1022. OriginalMinValue = res.Min(),
  1023. CreateTime = DateTime.Now,
  1024. StatStartTime = collection.OrderBy(i => i.Timestamp).Select(i => i.Timestamp).First(),
  1025. StatEndTime = collection.OrderBy(i => i.Timestamp).Select(i => i.Timestamp).Last(),
  1026. StatMaxValueFprCoefficient = statMaxValueFprCoefficient,
  1027. StatMinValueFprCoefficient = statMinValueFprCoefficient,
  1028. StatModeAvgFprCoefficient = StatModeAvgFprCoefficient,
  1029. Remark = string.Empty,
  1030. SerialTailNumber = serialNo.Substring(serialNo.Length - 2)
  1031. };
  1032. }
  1033. /// <summary>
  1034. /// 去除高频数据
  1035. /// </summary>
  1036. /// <param name="phr"></param>
  1037. /// <param name="highFreqSampleInterva"></param>
  1038. /// <returns></returns>
  1039. private static List<PregnancyHeartRateModel> GetNonFreqPregnancyHeartRate(List<PregnancyHeartRateModel> phr, int highFreqSampleInterval)
  1040. {
  1041. //phr = phr.OrderByDescending(i => i.LastUpdate).ToList();
  1042. //var result = new List<PregnancyHeartRateModel>();
  1043. //PregnancyHeartRateModel? previousItem = null;
  1044. //foreach (var item in phr)
  1045. //{
  1046. // if (previousItem != null)
  1047. // {
  1048. // var timeNextDiff =(previousItem!.LastUpdate - item.LastUpdate).TotalSeconds;
  1049. // if (timeNextDiff > highFreqSampleInterval)
  1050. // {
  1051. // result.Add(previousItem);
  1052. // }
  1053. // }
  1054. // previousItem = item;
  1055. //}
  1056. //// 添加上一个
  1057. //if (previousItem != null)
  1058. //{
  1059. // result.Add(previousItem);
  1060. //}
  1061. //return result;
  1062. #region 反向
  1063. var phr1 = phr.OrderByDescending(i => i.LastUpdate).ToList();
  1064. var result = new List<PregnancyHeartRateModel>();
  1065. PregnancyHeartRateModel? previousItem1 = null;
  1066. foreach (var item in phr1)
  1067. {
  1068. if (previousItem1 != null)
  1069. {
  1070. var timeNextDiff = (previousItem1!.LastUpdate - item.LastUpdate).TotalSeconds;
  1071. if (timeNextDiff > highFreqSampleInterval)
  1072. {
  1073. result.Add(previousItem1);
  1074. }
  1075. }
  1076. previousItem1 = item;
  1077. }
  1078. // 添加上一个
  1079. if (previousItem1 != null)
  1080. {
  1081. result.Add(previousItem1);
  1082. }
  1083. #endregion
  1084. #region 正向
  1085. var phr2 = phr.OrderByDescending(i => i.LastUpdate).ToList(); ;
  1086. var freqCollection = new List<PregnancyHeartRateModel>();
  1087. PregnancyHeartRateModel? previousItem = null;
  1088. foreach (var item in phr2)
  1089. {
  1090. if (previousItem != null)
  1091. {
  1092. var timeNextDiff = (previousItem!.LastUpdate - item.LastUpdate).TotalSeconds;
  1093. if (timeNextDiff <= highFreqSampleInterval)
  1094. {
  1095. freqCollection.Add(item);
  1096. }
  1097. }
  1098. previousItem = item;
  1099. }
  1100. //去除高频
  1101. foreach (var item in freqCollection)
  1102. {
  1103. phr2.Remove(item);
  1104. }
  1105. #endregion
  1106. // 交集
  1107. var commonElements = phr2.Intersect(result).ToList();
  1108. return commonElements;
  1109. }
  1110. /// <summary>
  1111. /// 获取高频数据
  1112. /// </summary>
  1113. /// <param name="phr"></param>
  1114. /// <param name="highFreqSampleInterval"></param>
  1115. /// <returns></returns>
  1116. private static List<PregnancyHeartRateModel> GetFreqPregnancyHeartRate(List<PregnancyHeartRateModel> phr, int highFreqSampleInterval)
  1117. {
  1118. phr = phr.OrderByDescending(i => i.LastUpdate).ToList();
  1119. var freqCollection = new List<PregnancyHeartRateModel>();
  1120. PregnancyHeartRateModel? previousItem = null;
  1121. foreach (var item in phr)
  1122. {
  1123. if (previousItem != null)
  1124. {
  1125. var timeNextDiff = (previousItem.LastUpdate - item.LastUpdate).TotalSeconds;
  1126. if (timeNextDiff <= highFreqSampleInterval)
  1127. {
  1128. freqCollection.Add(previousItem);
  1129. }
  1130. }
  1131. previousItem = item;
  1132. }
  1133. // 检查最后一条是否高频
  1134. if (previousItem != null && (phr.Last().LastUpdate - previousItem.LastUpdate).TotalSeconds <= highFreqSampleInterval)
  1135. {
  1136. freqCollection.Add(previousItem);
  1137. }
  1138. return freqCollection;
  1139. }
  1140. /// <summary>
  1141. /// 获取孕妇心率众数
  1142. /// </summary>
  1143. /// <param name="serialNo"></param>
  1144. /// <param name="days"></param>
  1145. /// <returns></returns>
  1146. //public async Task<int> GetPregnancyHeartRateModeAsync(string serialNo,int days=7)
  1147. //{
  1148. // var tableName = typeof(PregnancyHeartRateModel)
  1149. // .GetCustomAttribute<STableAttribute>()?
  1150. // .STableName;
  1151. // var daysAgo = DateTime.Now.AddDays(-days);
  1152. // var res = await _clientSqlSugar
  1153. // .Queryable<PregnancyHeartRateModel>()
  1154. // .AS(tableName)
  1155. // .Where(i=>i.SerialNumber.Equals(serialNo))
  1156. // .Where(i => i.Timestamp > daysAgo)
  1157. // //.OrderByDescending(i => i.PregnancyHeartRate)
  1158. // .Select(i =>i.PregnancyHeartRate)
  1159. // .ToListAsync();
  1160. // // 心率数据量必须30个以上才进行计算
  1161. // if (res.Count < 30) return 0;
  1162. // // 计算众数
  1163. // var mode = res.GroupBy(n => n)
  1164. // .OrderByDescending(g => g.Count())
  1165. // .First()
  1166. // .Key;
  1167. // Console.WriteLine("众数是: " + mode);
  1168. // // 如果有多个众数的情况
  1169. // var maxCount = res.GroupBy(n => n)
  1170. // .Max(g => g.Count());
  1171. // var modes = res.GroupBy(n => n)
  1172. // .Where(g => g.Count() == maxCount)
  1173. // .Select(g => g.Key)
  1174. // .ToList();
  1175. // // 多个众数,选择最接近平均数或中位数的众数
  1176. // if (modes.Count>1)
  1177. // {
  1178. // // 计算平均值
  1179. // double average = res.Average();
  1180. // Console.WriteLine("平均值是: " + average);
  1181. // // 计算中位数
  1182. // double median;
  1183. // int count = res.Count;
  1184. // var sortedRes = res.OrderBy(n => n).ToList();
  1185. // if (count % 2 == 0)
  1186. // {
  1187. // // 偶数个元素,取中间两个数的平均值
  1188. // median = (sortedRes[count / 2 - 1] + sortedRes[count / 2]) / 2.0;
  1189. // }
  1190. // else
  1191. // {
  1192. // // 奇数个元素,取中间的数
  1193. // median = sortedRes[count / 2];
  1194. // }
  1195. // Console.WriteLine("中位数是: " + median);
  1196. // // 找出最接近平均值的众数
  1197. // //var closestToAverage = modes.OrderBy(m => Math.Abs(m - average)).First();
  1198. // //Console.WriteLine("最接近平均值的众数是: " + closestToAverage);
  1199. // // 找出最接近中位数的众数
  1200. // var closestToMedian = modes.OrderBy(m => Math.Abs(m - median)).First();
  1201. // Console.WriteLine("最接近中位数的众数是: " + closestToMedian);
  1202. // mode = closestToMedian;
  1203. // }
  1204. // return mode;
  1205. //}
  1206. #endregion
  1207. }
  1208. }