vib.py 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324
  1. import os
  2. import time
  3. import traceback
  4. from datetime import datetime
  5. import pandas as pd
  6. import requests
  7. from rz_token import TokenManager, CONFIG, logger
  8. class CmsApiClient:
  9. """
  10. 统一CMS接口客户端
  11. 统一管理Token,内置:资产台账、振动指标数据 两个接口
  12. """
  13. def __init__(self, config=None):
  14. self.config = config or CONFIG
  15. # 全局唯一Token管理器
  16. self.token_manager = TokenManager(self.config)
  17. self.session = requests.Session()
  18. # 接口地址配置
  19. self.asset_api_url = self.config["base_url"] + "/api/openapiService/asset/v2/assetAccount"
  20. self.vib_temp_api_url = self.config["base_url"] + "/api/openapiService/sampleData/v2/vibTempData"
  21. def _get_common_headers(self):
  22. """统一获取带Token的请求头"""
  23. auth_header = self.token_manager.get_auth_header()
  24. return {
  25. "Authorization": auth_header,
  26. "Content-Type": "application/json; charset=utf-8"
  27. }
  28. # ===================== 1、获取设备资产台账(自动分页) =====================
  29. def get_asset_list(self, params=None, page_size=10000):
  30. headers = self._get_common_headers()
  31. query_params = {
  32. "skipCount": 0,
  33. "maxResultCount": page_size
  34. }
  35. if params:
  36. query_params.update(params)
  37. all_items = []
  38. page = 1
  39. while True:
  40. logger.info(f"获取资产 正在获取第 {page} 页数据...")
  41. try:
  42. logger.info(f"请求地址{self.asset_api_url},请求参数:{query_params},headers:{headers}")
  43. resp = self.session.post(
  44. url=self.asset_api_url,
  45. json=query_params,
  46. headers=headers,
  47. timeout=60
  48. )
  49. resp.raise_for_status()
  50. result = resp.json()
  51. except Exception as e:
  52. raise Exception(f"请求获取资产失败:{str(e)}")
  53. if not result.get("success"):
  54. raise Exception(f"获取资产接口失败:{result.get('message')}")
  55. data = result.get("data", {})
  56. items = data.get("items", [])
  57. total_count = data.get("totalCount", 0)
  58. if not items:
  59. break
  60. all_items.extend(items)
  61. logger.info(f"本页获取 {len(items)} 条,累计:{len(all_items)}/{total_count}")
  62. if len(all_items) >= total_count:
  63. break
  64. query_params["skipCount"] = len(all_items)
  65. page += 1
  66. time.sleep(0.2)
  67. logger.info(f"资产台账全部获取完成,共 {len(all_items)} 条")
  68. return all_items
  69. # ===================== 2、获取振动&指标(特征值)数据 =====================
  70. def get_vib_temp_data(self, params):
  71. begin_time = datetime.now()
  72. headers = self._get_common_headers()
  73. # logger.info(f"正在获取设备【{params.get('deviceCode')}】振动/指标数据...")
  74. # logger.info(f"请求地址{self.vib_temp_api_url},请求参数:{params}")
  75. try:
  76. resp = self.session.post(
  77. url=self.vib_temp_api_url,
  78. json=params,
  79. headers=headers,
  80. timeout=600
  81. )
  82. resp.raise_for_status()
  83. result = resp.json()
  84. except Exception as e:
  85. logger.error((traceback.format_exc()))
  86. # raise Exception(f"请求振动指标数据失败:{str(e)}")
  87. return {}
  88. if not result.get("success"):
  89. msg = f"振动指标接口失败:{result.get('message')}"
  90. logger.error(msg)
  91. # raise Exception(msg)
  92. return {}
  93. # logger.info(result.get("data"))
  94. logger.info(f"{str(params)},请求耗时:{datetime.now() - begin_time}")
  95. return result.get("data", {})
  96. def get_prev_day(date: str | None | datetime = None, prev_days: int = 1) -> str:
  97. """
  98. 获取上一天日期
  99. :param date: 日期字符串,格式为YYYY-MM-DD
  100. :param prev_days: 上几天,默认1天
  101. :return: 上几天日期字符串,格式为YYYY-MM-DD ,例如:2026-06-01
  102. """
  103. if not date:
  104. date = datetime.now().strftime("%Y-%m-%d")
  105. date = pd.to_datetime(date)
  106. prev_day = date - pd.Timedelta(days=prev_days)
  107. return prev_day.strftime("%Y-%m-%d")
  108. def get_prev_days(date: str | None | datetime = None, prev_days: int = 1) -> list:
  109. """
  110. 获取上一天日期
  111. :param date: 日期字符串,格式为YYYY-MM-DD
  112. :param prev_days: 上几天,默认1天
  113. :return: 上几天日期字符串,格式为YYYY-MM-DD ,例如:2026-06-01
  114. """
  115. if not date:
  116. date = datetime.now().strftime("%Y-%m-%d")
  117. date = pd.to_datetime(date)
  118. prev_datas = []
  119. if prev_days > 0:
  120. for i in range(prev_days, 0, -1):
  121. prev_datas.append((date - pd.Timedelta(days=i)).strftime("%Y-%m-%d"))
  122. else:
  123. for i in range(0, -prev_days + 1):
  124. prev_datas.append((date + pd.Timedelta(days=i)).strftime("%Y-%m-%d"))
  125. return prev_datas
  126. def complete_minute_data_fast(df: pd.DataFrame, query_time: str) -> pd.DataFrame:
  127. if df.empty:
  128. return df.copy()
  129. df = df.copy()
  130. df['localTime'] = pd.to_datetime(df['localTime']).dt.floor('min')
  131. group_cols = ['pointName', 'indexCode', 'departmentName']
  132. minute_range = pd.date_range(
  133. start=f'{query_time} 00:00:00',
  134. end=f'{query_time} 23:59:00',
  135. freq='min'
  136. )
  137. # 同一分组同一分钟只保留最后一条
  138. df = df.sort_values([*group_cols, 'localTime'])
  139. df = df.drop_duplicates([*group_cols, 'localTime'], keep='last')
  140. # 方式1:使用merge笛卡尔积,规避MultiIndex.from_product层级陷阱
  141. unique_keys = df[group_cols].drop_duplicates().reset_index(drop=True)
  142. time_df = pd.DataFrame({"localTime": minute_range})
  143. # 所有点位 × 所有分钟
  144. full_template = unique_keys.merge(time_df, how="cross")
  145. # 合并原始数据
  146. df_merge = full_template.merge(
  147. df[[*group_cols, "localTime", "value"]],
  148. on=[*group_cols, "localTime"],
  149. how="left"
  150. )
  151. # 分组内前向填充value
  152. df_merge["value"] = df_merge.groupby(group_cols, sort=False)["value"].ffill()
  153. # df_merge["value"] = df_merge.groupby(group_cols, sort=False)["value"].bfill()
  154. # 格式化时间字符串
  155. df_merge['localTime'] = df_merge['localTime'].dt.strftime('%Y-%m-%d %H:%M:00')
  156. return df_merge[['pointName', 'indexCode', 'departmentName', 'localTime', 'value']]
  157. def get_day_data(wind40, query_time):
  158. begin_time = datetime.now()
  159. index_codes = [[260, 261, 262], [263, 264, 265], [266, 276, 277]]
  160. department_map = {'发电机': '发电机', '齿轮箱': '齿轮箱', '行星': '齿轮箱', '主轴': '主轴'}
  161. result_jsons = list()
  162. id = wind40['id']
  163. code = wind40['code']
  164. name = wind40['name']
  165. fullName = wind40['fullName'].replace('/', '-')
  166. wind_turbine_type = wind40['wind_turbine_type']
  167. wind_farm_id = wind40['wind_farm_id']
  168. wind_turbine_id = wind40['wind_turbine_id']
  169. # childrens = wind40['children']
  170. # pointCodes = []
  171. # for child in childrens:
  172. # pointCodes.append(child['code'])
  173. for index_code in index_codes:
  174. vib_params = {
  175. "startTime": f"{query_time} 00:00:00",
  176. "endTime": f"{query_time} 23:59:59",
  177. "deviceCode": id,
  178. "types": [1],
  179. # "pointCodes": pointCodes,
  180. "vibFilter": {
  181. # 0:加速度,1:速度,2:位移
  182. "signalTypes": [1],
  183. },
  184. "indexFilter": {
  185. # 0:加速度,1:速度,2:位移
  186. "signalTypes": [1],
  187. "indexCodes": index_code
  188. }
  189. }
  190. vib_data = client.get_vib_temp_data(params=vib_params)
  191. if vib_data and vib_data['pointSampleDatas']:
  192. deviceId = vib_data['deviceId']
  193. deviceName = vib_data['deviceName']
  194. # deviceCode = vib_data['deviceCode']
  195. for data in vib_data['pointSampleDatas']:
  196. pointId = data['pointId']
  197. pointName = data['pointName']
  198. poinFullName = data['poinFullName']
  199. departmentName = None
  200. for k, v in department_map.items():
  201. if k in pointName:
  202. departmentName = v
  203. for dat in data['indexDatas']:
  204. measurementId = dat['measurementId']
  205. measurementName = dat['measurementName']
  206. length = dat['length']
  207. lowerFreq = dat['lowerFreq']
  208. upperFreq = dat['upperFreq']
  209. # signalType = dat['signalType']
  210. logger.info(
  211. f"{poinFullName},{str(index_code)} 获取到的数据数量为:{len(dat['datas'])}")
  212. if '16k 速度波形(10-1000)' == measurementName:
  213. for da in dat['datas']:
  214. indexCode = da['indexCode']
  215. da_time = da['time']
  216. value = da['value']
  217. # waveKey = da['waveKey']
  218. # condition = da['condition']
  219. # isErrorSignal = da['isErrorSignal']
  220. # unit = da['unit']
  221. result = {
  222. # 'deviceId': deviceId, 'deviceName': deviceName,
  223. # 'deviceCode': deviceCode,
  224. # 'pointId': pointId,
  225. 'pointName': pointName,
  226. # 'pointFullName': poinFullName,
  227. # 'measurementId': measurementId, 'measurementName': measurementName,
  228. # 'length': length, 'lowerFreq': lowerFreq, 'upperFreq': upperFreq,
  229. 'indexCode': indexCode, 'localTime': da_time,
  230. 'value': value, 'departmentName': departmentName
  231. }
  232. # logger.info(f'{now_index},{result}')
  233. result_jsons.append(result)
  234. else:
  235. logger.warn(f'{fullName} 无数据,休眠0.01秒')
  236. logger.info(f"{query_time},获取数据数量:{len(result_jsons)},耗时:{datetime.now() - begin_time}")
  237. if result_jsons:
  238. df = complete_minute_data_fast(pd.DataFrame(result_jsons), query_time)
  239. file_name = f'/data/wind_files/headquarter/cms/tmp/{query_time.replace('-','')}/{wind_turbine_type}/{wind_farm_id}/{wind_turbine_id}.parquet'
  240. os.makedirs(os.path.dirname(file_name), exist_ok=True)
  241. df.to_parquet(file_name, index=False)
  242. logger.info(f'{query_time},{fullName}保存成功,耗时:{datetime.now() - begin_time}')
  243. # ===================== 统一调用示例 =====================
  244. if __name__ == "__main__":
  245. app_begin_time = datetime.now()
  246. try:
  247. # 只初始化一次客户端,Token全局复用
  248. client = CmsApiClient()
  249. import json, multiprocessing, sys
  250. # file_name = 'wind_data_exists.json'
  251. file_suffix = sys.argv[1] if len(sys.argv) > 1 else '1'
  252. file_name = f'wind_data_with_scada_{file_suffix}.json'
  253. print(file_name)
  254. with open(file_name, 'r', encoding='utf-8') as f:
  255. datas = json.load(f)
  256. bool_has_result = 0
  257. time_datas = get_prev_days(prev_days=1)
  258. total = len(datas)
  259. for query_time in time_datas:
  260. with multiprocessing.Pool(6, maxtasksperchild=5) as pool:
  261. pool.starmap(get_day_data, [(v, query_time) for v in datas])
  262. # total = 1
  263. # get_day_data(get_prev_day(prev_days=1))
  264. logger.info(f"执行结束,总耗时:{datetime.now() - app_begin_time}")
  265. from combine_cms_data import main
  266. main()
  267. except Exception as e:
  268. logger.error(f"执行出错:{str(e)}")
  269. raise