collection_data.py 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302
  1. import hashlib
  2. import json
  3. import multiprocessing
  4. import os
  5. import sys
  6. import time
  7. from datetime import datetime
  8. from typing import Dict, List, Tuple
  9. from urllib.parse import urlparse, parse_qs
  10. import pandas as pd
  11. import requests
  12. class EnosAPIClient:
  13. TOKEN_CACHE_FILE = "enos_token_cache.json"
  14. def __init__(self, app_key: str, app_secret: str, apigw_address: str,
  15. org_id: str, is_edge: bool = True, cache_dir: str = None):
  16. self.app_key = app_key
  17. self.app_secret = app_secret
  18. self.apigw_address = apigw_address
  19. self.org_id = org_id
  20. self.protocol = "http" if is_edge else "https"
  21. self.cache_dir = cache_dir or os.path.join(os.getcwd(), ".enos_cache")
  22. os.makedirs(self.cache_dir, exist_ok=True)
  23. self.cache_file = os.path.join(self.cache_dir, self.TOKEN_CACHE_FILE)
  24. self.access_token = None
  25. self.token_expire_time = 0
  26. self._load_token_from_cache()
  27. def _sha256_lower(self, text: str) -> str:
  28. return hashlib.sha256(text.encode('utf-8')).hexdigest().lower()
  29. def _load_token_from_cache(self) -> bool:
  30. try:
  31. if not os.path.exists(self.cache_file):
  32. return False
  33. with open(self.cache_file, 'r', encoding='utf-8') as f:
  34. cache_data = json.load(f)
  35. cache_key = f"{self.app_key}_{self.apigw_address}"
  36. token_info = cache_data.get(cache_key)
  37. if token_info and time.time() < token_info.get('expire_time', 0) - 300:
  38. self.access_token = token_info.get('access_token')
  39. self.token_expire_time = token_info.get('expire_time')
  40. return True
  41. return False
  42. except:
  43. return False
  44. def _save_token_to_cache(self, access_token: str, expire_seconds: int):
  45. try:
  46. cache_data = {}
  47. if os.path.exists(self.cache_file):
  48. with open(self.cache_file, 'r', encoding='utf-8') as f:
  49. cache_data = json.load(f)
  50. cache_key = f"{self.app_key}_{self.apigw_address}"
  51. cache_data[cache_key] = {
  52. 'access_token': access_token,
  53. 'expire_time': time.time() + expire_seconds
  54. }
  55. with open(self.cache_file, 'w', encoding='utf-8') as f:
  56. json.dump(cache_data, f, indent=2)
  57. except:
  58. pass
  59. def get_access_token(self, force_refresh: bool = False) -> Tuple[bool, str]:
  60. if not force_refresh and self.access_token and time.time() < self.token_expire_time - 300:
  61. return True, "使用缓存Token"
  62. timestamp = int(time.time() * 1000)
  63. encryption = self._sha256_lower(f"{self.app_key}{timestamp}{self.app_secret}")
  64. url = f"{self.protocol}://{self.apigw_address}/apim-token-service/v2.0/token/get"
  65. try:
  66. resp = requests.post(url, json={
  67. "appKey": self.app_key,
  68. "encryption": encryption,
  69. "timestamp": timestamp
  70. }, timeout=10)
  71. result = resp.json()
  72. if result.get("status") == 0:
  73. token_data = result.get("data", {})
  74. self.access_token = token_data.get("accessToken")
  75. expire = token_data.get("expire", 7200)
  76. self.token_expire_time = time.time() + expire
  77. self._save_token_to_cache(self.access_token, expire)
  78. return True, "获取Token成功"
  79. return False, result.get("msg", "未知错误")
  80. except Exception as e:
  81. return False, str(e)
  82. def _generate_signature(self, url: str, body: str = "") -> Dict:
  83. parsed = urlparse(url)
  84. params = parse_qs(parsed.query)
  85. sorted_keys = sorted(params.keys())
  86. params_data = "".join([f"{k}{params[k][0]}" for k in sorted_keys if params[k]])
  87. if body:
  88. params_data += body
  89. timestamp = str(int(time.time() * 1000))
  90. sign_data = f"{self.access_token}{params_data}{timestamp}{self.app_secret}"
  91. apim_sign = self._sha256_lower(sign_data)
  92. return {
  93. "apim-accesstoken": self.access_token,
  94. "apim-signature": apim_sign,
  95. "apim-timestamp": timestamp,
  96. "Content-Type": "application/json; charset=utf-8"
  97. }
  98. def call_api(self, url: str, method: str = "GET", body: Dict = None) -> Tuple[bool, Dict]:
  99. if not self.access_token or time.time() >= self.token_expire_time:
  100. success, msg = self.get_access_token(force_refresh=True)
  101. if not success:
  102. return False, {"error": f"Token获取失败: {msg}"}
  103. body_str = json.dumps(body, ensure_ascii=False) if body else ""
  104. headers = self._generate_signature(url, body_str)
  105. try:
  106. if method.upper() == "GET":
  107. resp = requests.get(url, headers=headers, timeout=30)
  108. else:
  109. resp = requests.post(url, headers=headers, data=body_str.encode('utf-8'), timeout=30)
  110. result = resp.json()
  111. if result.get("code") == 0:
  112. return True, result
  113. return False, result
  114. except Exception as e:
  115. return False, {"error": str(e)}
  116. def split_list(lst: List, chunk_size: int) -> List[List]:
  117. """将列表按指定大小分组"""
  118. return [lst[i:i + chunk_size] for i in range(0, len(lst), chunk_size)]
  119. def query_historical_data(client: EnosAPIClient, mdm_id: str, point_ids: List[str],
  120. start_time: str, end_time: str, interval: str = "RAW", retry_times=0) -> Tuple[
  121. bool, List[Dict]]:
  122. """
  123. 查询历史测点数据(支持测点自动分批)
  124. """
  125. all_items = []
  126. point_groups = split_list(point_ids, 20) # 每20个测点一组
  127. print(f"共 {len(point_ids)} 个测点,分 {len(point_groups)} 批查询")
  128. for idx, group in enumerate(point_groups, 1):
  129. points_str = ",".join(group)
  130. url = (f"{client.protocol}://{client.apigw_address}/cds-timeseries-service/v1.0/tsdb-detail"
  131. f"?action=query"
  132. f"&orgId={client.org_id}"
  133. f"&mdmIds={mdm_id}"
  134. f"&pointIdsWithLogic={points_str}"
  135. f"&startTime={start_time}"
  136. f"&endTime={end_time}"
  137. f"&interval={interval}")
  138. success, result = client.call_api(url, "GET")
  139. if success:
  140. items = result.get("data", {}).get("items", [])
  141. all_items.extend(items)
  142. print(f"批次 {idx}/{len(point_groups)}:获取 {len(items)} 条记录")
  143. else:
  144. if retry_times <= 5:
  145. retry_times = retry_times + 1
  146. error_msg = result.get("error", result.get("msg", "未知错误"))
  147. print(f"批次 {idx}/{len(point_groups)}: 查询失败 - {error_msg},次数: {retry_times}")
  148. query_historical_data(client, mdm_id, point_ids, start_time, end_time, interval, retry_times)
  149. time.sleep(0.01)
  150. else:
  151. print(f"风机:{mdm_id}, {len(point_groups)}, 查询{retry_times}仍失败")
  152. return False, []
  153. return True, all_items
  154. def get_prev_day(date: str | None | datetime = None, prev_days: int = 1) -> str:
  155. """
  156. 获取上一天日期
  157. :param date: 日期字符串,格式为YYYY-MM-DD
  158. :param prev_days: 上几天,默认1天
  159. :return: 上几天日期字符串,格式为YYYY-MM-DD ,例如:2026-06-01
  160. """
  161. if not date:
  162. date = datetime.now().strftime("%Y-%m-%d")
  163. date = pd.to_datetime(date)
  164. prev_day = date - pd.Timedelta(days=prev_days)
  165. return prev_day.strftime("%Y-%m-%d")
  166. def get_prev_days(date: str | None | datetime = None, prev_days: int = 1) -> list:
  167. """
  168. 获取上一天日期
  169. :param date: 日期字符串,格式为YYYY-MM-DD
  170. :param prev_days: 上几天,默认1天
  171. :return: 上几天日期字符串,格式为YYYY-MM-DD ,例如:2026-06-01
  172. """
  173. if not date:
  174. date = datetime.now().strftime("%Y-%m-%d")
  175. date = pd.to_datetime(date)
  176. prev_datas = []
  177. if prev_days > 0:
  178. for i in range(prev_days, 0, -1):
  179. prev_datas.append((date - pd.Timedelta(days=i)).strftime("%Y-%m-%d"))
  180. else:
  181. for i in range(0, -prev_days + 1):
  182. prev_datas.append((date + pd.Timedelta(days=i)).strftime("%Y-%m-%d"))
  183. return prev_datas
  184. def query_by_day(client: EnosAPIClient, mdm_id: str, point_ids: List[str], query_date: str) -> None:
  185. """
  186. 按天分批查询数据
  187. """
  188. json_datas = {}
  189. day_start = f"{query_date} 00:00:00"
  190. day_end = f"{query_date} 23:59:59"
  191. print(f"查询日期: {query_date}")
  192. success, items = query_historical_data(
  193. client, mdm_id, point_ids, day_start, day_end, "RAW"
  194. )
  195. if success:
  196. for item in items:
  197. localtime = item.get("localtime")
  198. if not localtime in json_datas.keys():
  199. json_datas[localtime] = {}
  200. json_datas[localtime].update(item)
  201. print(f"当天总计: {len(items)} 条记录")
  202. else:
  203. print(f"查询失败: {items}")
  204. # 保存所有数据
  205. if json_datas:
  206. df = pd.DataFrame().from_dict(json_datas.values())
  207. filename = f"{mdm_id}_{query_date}.parquet"
  208. df.to_parquet(filename, index=False)
  209. print(f"总共获取 {df.shape} 条记录,已保存到 {filename}")
  210. else:
  211. print(" 未获取到任何数据")
  212. def get_wind_farm(read_db):
  213. df_wind_farm = pd.read_csv('wind_farm.csv', encoding='utf-8')
  214. df_wind_farm_db = df_wind_farm[df_wind_farm['store_id'] == read_db]
  215. return df_wind_farm_db
  216. def get_wind_turbine():
  217. pass
  218. # ===================== 使用示例 =====================
  219. def main(read_db):
  220. # 配置参数(请替换为实际值)
  221. APP_KEY = "a6af2233-227e-4684-ab59-9e5f74e5718c"
  222. APP_SECRET = "56f24801-0429-486f-b4db-ba61a33a101c"
  223. API_GW_ADDRESS = "ag-cdt1.eniot.io"
  224. ORG_ID = "o16021383932361"
  225. # 设备ID列表
  226. MDM_IDS = ["03puSZ72"]
  227. # 测点ID列表(100多个测点,会自动按20个一组分批)
  228. POINT_IDS = "last(WCNV.DECCNVAI005),last(WCNV.DECCNVAI008),last(WCNV.DECCNVAI006),last(WCNV.DECCNVAI009),last(WCNV.DECCNVAI007),last(WCNV.DECCNVAI010),last(WTRM.TrmTmpShfBrg),first(WYAW.DECYAWDI004),last(WYAW.DECYAWAI001),first(WTUR.TurbineSts),first(WTUR.TurbineSts_Map),last(WGEN.GenSpd),last(WGEN.TemGenNonDE),last(WGEN.TemGenDriEnd),last(WROT.TemAxis1Ctrl),last(WROT.TemAxis2Ctrl),last(WROT.TemAxis3Ctrl),last(WROT.TemB1Mot),last(WROT.TemB2Mot),last(WROT.TemB3Mot),last(WROT.CurBlade1Motor),last(WROT.CurBlade2Motor),last(WROT.CurBlade3Motor),last(WROT.DECROTAI004),last(WROT.DECROTAI005),last(WROT.DECROTAI006),last(WROT.Blade1Speed),last(WROT.Blade2Speed),last(WROT.Blade3Speed),last(WCNV.DECCNVAI012),last(WCNV.DECCNVAI011),last(WTOW.TemTower),last(WTUR.DECTURAI016),last(WTUR.DECTURAI015),first(WTUR.DECTURDI030),last(WGEN.TemGenStaU),last(WGEN.TemGenStaV),last(WGEN.TemGenStaW),last(WGEN.Torque),last(WYAW.TotalTwist),first(WTUR.AIStatusCode),last(WGEN.GenReactivePW),last(WGEN.GenActivePW),last(WTUR.DECTURAI003),last(WYAW.NacellePosition),last(WNAC.TemNacelle),last(WNAC.TemOut),last(WNAC.TemNacelleCab),last(WROT.Blade1Position),last(WROT.Blade2Position),last(WROT.Blade3Position),last(WROT.Blade1Setpoint),last(WROT.Blade2Setpoint),last(WROT.Blade3Setpoint),last(WTRM.PowerStoreTRBS),last(WNAC.DECNACAI004),last(WNAC.TheoryActivePW),last(WCNV.VolConL1),last(WCNV.CurConL1),last(WCNV.VolConL2),last(WCNV.CurConL2),last(WCNV.VolConL3),last(WCNV.CurConL3),last(WGEN.TorqueSetpoint),last(WNAC.DECNACAI003),last(WCNV.GridFreq),first(WYAW.YawCCWSts),first(WTUR.CCTurbineSts),first(WYAW.YawCWSts),last(WNAC.WindDirection),first(WTUR.GenState),last(WTRM.RotorSpd),last(WNAC.WindSpeed),last(WTRM.TemGeaMSND),last(WTRM.TemGeaMSDE),last(WTRM.GBoxOilPmpP),last(WTRM.TemGeaOil),last(WTRM.GBoxSpd)".split(
  229. ",")
  230. POINT_IDS = POINT_IDS[0:15]
  231. # 初始化客户端
  232. client = EnosAPIClient(APP_KEY, APP_SECRET, API_GW_ADDRESS, ORG_ID)
  233. # 获取Token
  234. success, msg = client.get_access_token()
  235. if not success:
  236. print(f"error: {msg}")
  237. return
  238. print(f"success: {msg}")
  239. for query_date in get_prev_days(prev_days=1):
  240. with multiprocessing.Pool(6, maxtasksperchild=5) as pool:
  241. # 查询数据
  242. pool.starmap(query_by_day, [(client, mdm_id, POINT_IDS, query_date) for mdm_id in MDM_IDS])
  243. if __name__ == "__main__":
  244. read_db = 1
  245. args = sys.argv
  246. if len(args) > 1:
  247. read_db = int(args[1])
  248. print(read_db)
  249. main(read_db)