您最多选择25个主题 主题必须以字母或数字开头,可以包含连字符 (-),并且长度不得超过35个字符

TDengineService.cs 17KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427
  1. using HealthMonitor.Common;
  2. using HealthMonitor.Common.helper;
  3. using HealthMonitor.Model.Config;
  4. using HealthMonitor.Service.Biz.db.Dto;
  5. using Microsoft.EntityFrameworkCore.Metadata.Internal;
  6. using Microsoft.Extensions.Logging;
  7. using Microsoft.Extensions.Options;
  8. using Newtonsoft.Json;
  9. using System;
  10. using System.Collections.Generic;
  11. using System.Linq;
  12. using System.Text;
  13. using System.Threading.Tasks;
  14. using TDengineDriver;
  15. using TDengineDriver.Impl;
  16. namespace HealthMonitor.Service.Biz.db
  17. {
  18. public class TDengineService
  19. {
  20. private readonly ILogger<TDengineService> _logger;
  21. private readonly HttpHelper _httpHelper=default!;
  22. private readonly TDengineServiceConfig _configTDengineService;
  23. public TDengineService(ILogger<TDengineService> logger,
  24. IOptions<TDengineServiceConfig> configTDengineService,
  25. HttpHelper httpHelper
  26. )
  27. {
  28. _logger = logger;
  29. _configTDengineService = configTDengineService.Value;
  30. _httpHelper = httpHelper;
  31. }
  32. public IntPtr Connection()
  33. {
  34. string host = _configTDengineService.Host;
  35. string user = _configTDengineService.UserName;
  36. string db = _configTDengineService.DB;
  37. short port = _configTDengineService.Port;
  38. string password = _configTDengineService.Password;
  39. IntPtr conn = TDengine.Connect(host, user, password, db, port);
  40. //string host = "172.16.255.180";
  41. //short port = 6030;
  42. //string user = "root";
  43. //string password = "taosdata";
  44. //string db = "health_monitor";
  45. //IntPtr conn = TDengine.Connect(host, user, password, db, port);
  46. // Check if get connection success
  47. if (conn == IntPtr.Zero)
  48. {
  49. Console.WriteLine("Connect to TDengine failed");
  50. }
  51. else
  52. {
  53. Console.WriteLine("Connect to TDengine success");
  54. }
  55. return conn;
  56. }
  57. public void ExecuteSQL(IntPtr conn, string sql)
  58. {
  59. IntPtr res = TDengine.Query(conn, sql);
  60. // Check if query success
  61. if ((res == IntPtr.Zero) || (TDengine.ErrorNo(res) != 0))
  62. {
  63. Console.Write(sql + " failure, ");
  64. // Get error message while Res is a not null pointer.
  65. if (res != IntPtr.Zero)
  66. {
  67. Console.Write("reason:" + TDengine.Error(res));
  68. }
  69. }
  70. else
  71. {
  72. Console.Write(sql + " success, {0} rows affected", TDengine.AffectRows(res));
  73. //... do something with res ...
  74. // Important: need to free result to avoid memory leak.
  75. TDengine.FreeResult(res);
  76. }
  77. }
  78. public void ExecuteQuerySQL(IntPtr conn, string sql)
  79. {
  80. IntPtr res = TDengine.Query(conn, sql);
  81. // Check if query success
  82. if ((res == IntPtr.Zero) || (TDengine.ErrorNo(res) != 0))
  83. {
  84. Console.Write(sql + " failure, ");
  85. // Get error message while Res is a not null pointer.
  86. if (res != IntPtr.Zero)
  87. {
  88. Console.Write("reason:" + TDengine.Error(res));
  89. }
  90. }
  91. else
  92. {
  93. Console.Write(sql + " success, {0} rows affected", TDengine.AffectRows(res));
  94. //... do something with res ...
  95. List<TDengineDriver.TDengineMeta> resMeta = LibTaos.GetMeta(res);
  96. List<Object> resData = LibTaos.GetData(res);
  97. foreach (var meta in resMeta)
  98. {
  99. Console.Write($"\t|{meta.name} {meta.TypeName()} ({meta.size})\t|");
  100. }
  101. for (int i = 0; i < resData.Count; i++)
  102. {
  103. Console.Write($"|{resData[i].ToString()} \t");
  104. //Console.WriteLine("{0},{1},{2}", i, resMeta.Count, (i) % resMeta.Count);
  105. if (((i + 1) % resMeta.Count == 0))
  106. {
  107. Console.WriteLine("");
  108. }
  109. }
  110. // Important: need to free result to avoid memory leak.
  111. TDengine.FreeResult(res);
  112. }
  113. }
  114. public void CheckRes(IntPtr conn, IntPtr res, String errorMsg)
  115. {
  116. if (TDengine.ErrorNo(res) != 0)
  117. {
  118. throw new Exception($"{errorMsg} since: {TDengine.Error(res)}");
  119. }
  120. }
  121. //public void ExecuteInsertSQL(IntPtr conn, string sql)
  122. //{
  123. // try
  124. // {
  125. // //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) " +
  126. // // "d1002 USING power.meters TAGS('California.SanFrancisco', 3) VALUES('2018-10-03 14:38:16.650', 10.30000, 218, 0.25000) " +
  127. // // "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) " +
  128. // // "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)";
  129. // IntPtr res = TDengine.Query(conn, sql);
  130. // CheckRes(conn, res, "failed to insert data");
  131. // int affectedRows = TDengine.AffectRows(res);
  132. // Console.WriteLine("affectedRows " + affectedRows);
  133. // TDengine.FreeResult(res);
  134. // }
  135. // finally
  136. // {
  137. // TDengine.Close(conn);
  138. // }
  139. //}
  140. public void ExecuteInsertSQL(string sql)
  141. {
  142. var conn = Connection();
  143. try
  144. {
  145. //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) " +
  146. // "d1002 USING power.meters TAGS('California.SanFrancisco', 3) VALUES('2018-10-03 14:38:16.650', 10.30000, 218, 0.25000) " +
  147. // "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) " +
  148. // "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)";
  149. IntPtr res = TDengine.Query(conn, sql);
  150. CheckRes(conn, res, "failed to insert data");
  151. int affectedRows = TDengine.AffectRows(res);
  152. Console.WriteLine("affectedRows " + affectedRows);
  153. TDengine.FreeResult(res);
  154. }
  155. finally
  156. {
  157. TDengine.Close(conn);
  158. }
  159. }
  160. #region TDengine.Connector async query
  161. public void QueryCallback(IntPtr param, IntPtr taosRes, int code)
  162. {
  163. if (code == 0 && taosRes != IntPtr.Zero)
  164. {
  165. FetchRawBlockAsyncCallback fetchRowAsyncCallback = new FetchRawBlockAsyncCallback(FetchRawBlockCallback);
  166. TDengine.FetchRawBlockAsync(taosRes, fetchRowAsyncCallback, param);
  167. }
  168. else
  169. {
  170. Console.WriteLine($"async query data failed, failed code {code}");
  171. }
  172. }
  173. // Iteratively call this interface until "numOfRows" is no greater than 0.
  174. public void FetchRawBlockCallback(IntPtr param, IntPtr taosRes, int numOfRows)
  175. {
  176. if (numOfRows > 0)
  177. {
  178. Console.WriteLine($"{numOfRows} rows async retrieved");
  179. IntPtr pdata = TDengine.GetRawBlock(taosRes);
  180. List<TDengineMeta> metaList = TDengine.FetchFields(taosRes);
  181. List<object> dataList = LibTaos.ReadRawBlock(pdata, metaList, numOfRows);
  182. for (int i = 0; i < metaList.Count; i++)
  183. {
  184. Console.Write("{0} {1}({2}) \t|", metaList[i].name, metaList[i].type, metaList[i].size);
  185. }
  186. Console.WriteLine();
  187. for (int i = 0; i < dataList.Count; i++)
  188. {
  189. if (i != 0 && i % metaList.Count == 0)
  190. {
  191. Console.WriteLine("{0}\t|", dataList[i]);
  192. }
  193. Console.Write("{0}\t|", dataList[i]);
  194. }
  195. TDengine.FetchRawBlockAsync(taosRes, FetchRawBlockCallback, param);
  196. }
  197. else
  198. {
  199. if (numOfRows == 0)
  200. {
  201. Console.WriteLine("async retrieve complete.");
  202. }
  203. else
  204. {
  205. Console.WriteLine($"FetchRawBlockCallback callback error, error code {numOfRows}");
  206. }
  207. TDengine.FreeResult(taosRes);
  208. }
  209. }
  210. public void ExecuteQueryAsync(string sql)
  211. {
  212. var conn = Connection();
  213. QueryAsyncCallback queryAsyncCallback = new QueryAsyncCallback(QueryCallback);
  214. TDengine.QueryAsync(conn, sql, queryAsyncCallback, IntPtr.Zero);
  215. }
  216. //public void ExecuteQuery(string sql)
  217. //{
  218. // var conn = Connection();
  219. // QueryAsyncCallback queryAsyncCallback = new QueryAsyncCallback(QueryCallback);
  220. // TDengine.QueryAsync(conn, sql, queryAsyncCallback, IntPtr.Zero);
  221. //}
  222. public Aggregate GetAggregateValue(string field, string tbName, string? condition)
  223. {
  224. List<int> data = new();
  225. var sql = $"SELECT MAX({field}), MIN({field}) FROM {_configTDengineService.DB}.{tbName} WHERE {condition}";
  226. var conn = Connection();
  227. try
  228. {
  229. IntPtr res = TDengine.Query(conn, sql);
  230. // Check if query success
  231. if ((res == IntPtr.Zero) || (TDengine.ErrorNo(res) != 0))
  232. {
  233. Console.Write(sql + " failure, ");
  234. // Get error message while Res is a not null pointer.
  235. if (res != IntPtr.Zero)
  236. {
  237. Console.Write("reason:" + TDengine.Error(res));
  238. }
  239. }
  240. else
  241. {
  242. Console.Write(sql + " success, {0} rows affected", TDengine.AffectRows(res));
  243. //... do something with res ...
  244. List<TDengineMeta> resMeta = LibTaos.GetMeta(res);
  245. List<object> resData = LibTaos.GetData(res);
  246. foreach (var meta in resMeta)
  247. {
  248. Console.Write($"\t|{meta.name} {meta.TypeName()} ({meta.size})\t|");
  249. }
  250. resData.ForEach(x => data.Add(SafeType.SafeInt(x)));
  251. // Important: need to free result to avoid memory leak.
  252. TDengine.FreeResult(res);
  253. }
  254. }
  255. finally
  256. {
  257. TDengine.Close(conn);
  258. }
  259. return new Aggregate
  260. {
  261. Max = data.Count.Equals(0) ? 0 : data[0],
  262. Min = data.Count.Equals(0) ? 0 : data[1],
  263. };
  264. }
  265. public int GetAvgExceptMaxMinValue(string field, string tbName, string? condition)
  266. {
  267. List<int> data = new();
  268. var sql = $"SELECT MAX({field}), MIN({field}) FROM {_configTDengineService.DB}.{tbName} WHERE {condition}";
  269. var aggregate= GetAggregateValue(field, tbName, condition);
  270. var sqlAvg = $"SELECT AVG({field}) FROM {_configTDengineService.DB}.{tbName} WHERE {condition} AND {field} < {aggregate.Max} and {field} > {aggregate.Min}";
  271. var conn = Connection();
  272. try
  273. {
  274. IntPtr res = TDengine.Query(conn, sqlAvg);
  275. // Check if query success
  276. if ((res == IntPtr.Zero) || (TDengine.ErrorNo(res) != 0))
  277. {
  278. Console.Write(sqlAvg + " failure, ");
  279. // Get error message while Res is a not null pointer.
  280. if (res != IntPtr.Zero)
  281. {
  282. Console.Write("reason:" + TDengine.Error(res));
  283. }
  284. }
  285. else
  286. {
  287. Console.Write(sqlAvg + " success, {0} rows affected", TDengine.AffectRows(res));
  288. //... do something with res ...
  289. List<TDengineMeta> resMeta = LibTaos.GetMeta(res);
  290. List<object> resData = LibTaos.GetData(res);
  291. foreach (var meta in resMeta)
  292. {
  293. Console.Write($"\t|{meta.name} {meta.TypeName()} ({meta.size})\t|");
  294. }
  295. resData.ForEach(x => data.Add(SafeType.SafeInt(x)));
  296. // Important: need to free result to avoid memory leak.
  297. TDengine.FreeResult(res);
  298. }
  299. }
  300. finally
  301. {
  302. TDengine.Close(conn);
  303. }
  304. return data.Count.Equals(0) ? 0 : data[0];
  305. }
  306. #endregion
  307. #region RestAPI
  308. public async Task<bool> GernalRestSql(string sql)
  309. {
  310. //"http://{server}:{port}/rest/sql/{db}"
  311. var url = $"http://{_configTDengineService.Host}:{_configTDengineService.RestPort}/rest/sql/{_configTDengineService.DB}";
  312. List<KeyValuePair<string, string>> headers = new()
  313. {
  314. new KeyValuePair<string, string>("Authorization", "Basic " + _configTDengineService.Token)
  315. };
  316. var result = await _httpHelper.HttpToPostAsync(url, sql, headers).ConfigureAwait(false);
  317. var res = JsonConvert.DeserializeObject<TDengineRestResBase>(result!);
  318. if (result != null)
  319. {
  320. if (res?.Code == 0)
  321. {
  322. _logger.LogInformation($"{nameof(GernalRestSql)},SQL 语句执行成功|{sql}");
  323. return true;
  324. }
  325. else
  326. {
  327. _logger.LogWarning($"{nameof(GernalRestSql)},SQL 语句执行失败||{sql}");
  328. return false;
  329. }
  330. }
  331. else
  332. {
  333. _logger.LogError($"{nameof(GernalRestSql)},TDengine 服务器IP:{_configTDengineService.Host} 错误,请联系运维人员");
  334. return false;
  335. }
  336. //return res.Code==0;
  337. }
  338. public async Task<string?> GernalRestSqlResTextAsync(string sql)
  339. {
  340. var url = $"http://{_configTDengineService.Host}:{_configTDengineService.RestPort}/rest/sql/{_configTDengineService.DB}";
  341. List<KeyValuePair<string, string>> headers = new()
  342. {
  343. new KeyValuePair<string, string>("Authorization", "Basic " + _configTDengineService.Token)
  344. };
  345. var result = await _httpHelper.HttpToPostAsync(url, sql, headers).ConfigureAwait(false);
  346. return result;
  347. }
  348. public async Task<Aggregate> GetAggregateValueAsync(string field,string tbName,string? condition)
  349. {
  350. var sql = $"SELECT MAX({field}), MIN({field}) FROM {_configTDengineService.DB}.{tbName} WHERE {condition}";
  351. var result = await GernalRestSqlResTextAsync(sql);
  352. var res = JsonConvert.DeserializeObject<Aggregate>(result!);
  353. List<dynamic> data = res?.Data!;
  354. return new Aggregate
  355. {
  356. Max = data.Count.Equals(0) ? 0 : data[0][0],
  357. Min = data.Count.Equals(0) ? 0 : data[0][1],
  358. };
  359. }
  360. public async Task<int> GetAvgExceptMaxMinValueAsync(string field, string tbName, string? condition)
  361. {
  362. var sql = $"SELECT MAX({field}), MIN({field}) FROM {_configTDengineService.DB}.{tbName} WHERE {condition}";
  363. var result = await GernalRestSqlResTextAsync(sql);
  364. var res = JsonConvert.DeserializeObject<TDengineRestResBase>(result!);
  365. List<dynamic> data = res?.Data!;
  366. 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])}";
  367. result = await GernalRestSqlResTextAsync(sqlAvg);
  368. res = JsonConvert.DeserializeObject<TDengineRestResBase>(result!);
  369. data = res?.Data!;
  370. return data.Count.Equals(0)?0:(int)data[0][0];
  371. }
  372. #endregion
  373. }
  374. }