| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324 |
- import os
- import time
- import traceback
- from datetime import datetime
- import pandas as pd
- import requests
- from rz_token import TokenManager, CONFIG, logger
- class CmsApiClient:
- """
- 统一CMS接口客户端
- 统一管理Token,内置:资产台账、振动指标数据 两个接口
- """
- def __init__(self, config=None):
- self.config = config or CONFIG
- # 全局唯一Token管理器
- self.token_manager = TokenManager(self.config)
- self.session = requests.Session()
- # 接口地址配置
- self.asset_api_url = self.config["base_url"] + "/api/openapiService/asset/v2/assetAccount"
- self.vib_temp_api_url = self.config["base_url"] + "/api/openapiService/sampleData/v2/vibTempData"
- def _get_common_headers(self):
- """统一获取带Token的请求头"""
- auth_header = self.token_manager.get_auth_header()
- return {
- "Authorization": auth_header,
- "Content-Type": "application/json; charset=utf-8"
- }
- # ===================== 1、获取设备资产台账(自动分页) =====================
- def get_asset_list(self, params=None, page_size=10000):
- headers = self._get_common_headers()
- query_params = {
- "skipCount": 0,
- "maxResultCount": page_size
- }
- if params:
- query_params.update(params)
- all_items = []
- page = 1
- while True:
- logger.info(f"获取资产 正在获取第 {page} 页数据...")
- try:
- logger.info(f"请求地址{self.asset_api_url},请求参数:{query_params},headers:{headers}")
- resp = self.session.post(
- url=self.asset_api_url,
- json=query_params,
- headers=headers,
- timeout=60
- )
- resp.raise_for_status()
- result = resp.json()
- except Exception as e:
- raise Exception(f"请求获取资产失败:{str(e)}")
- if not result.get("success"):
- raise Exception(f"获取资产接口失败:{result.get('message')}")
- data = result.get("data", {})
- items = data.get("items", [])
- total_count = data.get("totalCount", 0)
- if not items:
- break
- all_items.extend(items)
- logger.info(f"本页获取 {len(items)} 条,累计:{len(all_items)}/{total_count}")
- if len(all_items) >= total_count:
- break
- query_params["skipCount"] = len(all_items)
- page += 1
- time.sleep(0.2)
- logger.info(f"资产台账全部获取完成,共 {len(all_items)} 条")
- return all_items
- # ===================== 2、获取振动&指标(特征值)数据 =====================
- def get_vib_temp_data(self, params):
- begin_time = datetime.now()
- headers = self._get_common_headers()
- # logger.info(f"正在获取设备【{params.get('deviceCode')}】振动/指标数据...")
- # logger.info(f"请求地址{self.vib_temp_api_url},请求参数:{params}")
- try:
- resp = self.session.post(
- url=self.vib_temp_api_url,
- json=params,
- headers=headers,
- timeout=600
- )
- resp.raise_for_status()
- result = resp.json()
- except Exception as e:
- logger.error((traceback.format_exc()))
- # raise Exception(f"请求振动指标数据失败:{str(e)}")
- return {}
- if not result.get("success"):
- msg = f"振动指标接口失败:{result.get('message')}"
- logger.error(msg)
- # raise Exception(msg)
- return {}
- # logger.info(result.get("data"))
- logger.info(f"{str(params)},请求耗时:{datetime.now() - begin_time}")
- return result.get("data", {})
- def get_prev_day(date: str | None | datetime = None, prev_days: int = 1) -> str:
- """
- 获取上一天日期
- :param date: 日期字符串,格式为YYYY-MM-DD
- :param prev_days: 上几天,默认1天
- :return: 上几天日期字符串,格式为YYYY-MM-DD ,例如:2026-06-01
- """
- if not date:
- date = datetime.now().strftime("%Y-%m-%d")
- date = pd.to_datetime(date)
- prev_day = date - pd.Timedelta(days=prev_days)
- return prev_day.strftime("%Y-%m-%d")
- def get_prev_days(date: str | None | datetime = None, prev_days: int = 1) -> list:
- """
- 获取上一天日期
- :param date: 日期字符串,格式为YYYY-MM-DD
- :param prev_days: 上几天,默认1天
- :return: 上几天日期字符串,格式为YYYY-MM-DD ,例如:2026-06-01
- """
- if not date:
- date = datetime.now().strftime("%Y-%m-%d")
- date = pd.to_datetime(date)
- prev_datas = []
- if prev_days > 0:
- for i in range(prev_days, 0, -1):
- prev_datas.append((date - pd.Timedelta(days=i)).strftime("%Y-%m-%d"))
- else:
- for i in range(0, -prev_days + 1):
- prev_datas.append((date + pd.Timedelta(days=i)).strftime("%Y-%m-%d"))
- return prev_datas
- def complete_minute_data_fast(df: pd.DataFrame, query_time: str) -> pd.DataFrame:
- if df.empty:
- return df.copy()
- df = df.copy()
- df['localTime'] = pd.to_datetime(df['localTime']).dt.floor('min')
- group_cols = ['pointName', 'indexCode', 'departmentName']
- minute_range = pd.date_range(
- start=f'{query_time} 00:00:00',
- end=f'{query_time} 23:59:00',
- freq='min'
- )
- # 同一分组同一分钟只保留最后一条
- df = df.sort_values([*group_cols, 'localTime'])
- df = df.drop_duplicates([*group_cols, 'localTime'], keep='last')
- # 方式1:使用merge笛卡尔积,规避MultiIndex.from_product层级陷阱
- unique_keys = df[group_cols].drop_duplicates().reset_index(drop=True)
- time_df = pd.DataFrame({"localTime": minute_range})
- # 所有点位 × 所有分钟
- full_template = unique_keys.merge(time_df, how="cross")
- # 合并原始数据
- df_merge = full_template.merge(
- df[[*group_cols, "localTime", "value"]],
- on=[*group_cols, "localTime"],
- how="left"
- )
- # 分组内前向填充value
- df_merge["value"] = df_merge.groupby(group_cols, sort=False)["value"].ffill()
- # df_merge["value"] = df_merge.groupby(group_cols, sort=False)["value"].bfill()
- # 格式化时间字符串
- df_merge['localTime'] = df_merge['localTime'].dt.strftime('%Y-%m-%d %H:%M:00')
- return df_merge[['pointName', 'indexCode', 'departmentName', 'localTime', 'value']]
- def get_day_data(wind40, query_time):
- begin_time = datetime.now()
- index_codes = [[260, 261, 262], [263, 264, 265], [266, 276, 277]]
- department_map = {'发电机': '发电机', '齿轮箱': '齿轮箱', '行星': '齿轮箱', '主轴': '主轴'}
- result_jsons = list()
- id = wind40['id']
- code = wind40['code']
- name = wind40['name']
- fullName = wind40['fullName'].replace('/', '-')
- wind_turbine_type = wind40['wind_turbine_type']
- wind_farm_id = wind40['wind_farm_id']
- wind_turbine_id = wind40['wind_turbine_id']
-
- # childrens = wind40['children']
- # pointCodes = []
- # for child in childrens:
- # pointCodes.append(child['code'])
- for index_code in index_codes:
- vib_params = {
- "startTime": f"{query_time} 00:00:00",
- "endTime": f"{query_time} 23:59:59",
- "deviceCode": id,
- "types": [1],
- # "pointCodes": pointCodes,
- "vibFilter": {
- # 0:加速度,1:速度,2:位移
- "signalTypes": [1],
- },
- "indexFilter": {
- # 0:加速度,1:速度,2:位移
- "signalTypes": [1],
- "indexCodes": index_code
- }
- }
- vib_data = client.get_vib_temp_data(params=vib_params)
- if vib_data and vib_data['pointSampleDatas']:
- deviceId = vib_data['deviceId']
- deviceName = vib_data['deviceName']
- # deviceCode = vib_data['deviceCode']
- for data in vib_data['pointSampleDatas']:
- pointId = data['pointId']
- pointName = data['pointName']
- poinFullName = data['poinFullName']
- departmentName = None
- for k, v in department_map.items():
- if k in pointName:
- departmentName = v
- for dat in data['indexDatas']:
- measurementId = dat['measurementId']
- measurementName = dat['measurementName']
- length = dat['length']
- lowerFreq = dat['lowerFreq']
- upperFreq = dat['upperFreq']
- # signalType = dat['signalType']
- logger.info(
- f"{poinFullName},{str(index_code)} 获取到的数据数量为:{len(dat['datas'])}")
- if '16k 速度波形(10-1000)' == measurementName:
- for da in dat['datas']:
- indexCode = da['indexCode']
- da_time = da['time']
- value = da['value']
- # waveKey = da['waveKey']
- # condition = da['condition']
- # isErrorSignal = da['isErrorSignal']
- # unit = da['unit']
- result = {
- # 'deviceId': deviceId, 'deviceName': deviceName,
- # 'deviceCode': deviceCode,
- # 'pointId': pointId,
- 'pointName': pointName,
- # 'pointFullName': poinFullName,
- # 'measurementId': measurementId, 'measurementName': measurementName,
- # 'length': length, 'lowerFreq': lowerFreq, 'upperFreq': upperFreq,
- 'indexCode': indexCode, 'localTime': da_time,
- 'value': value, 'departmentName': departmentName
- }
- # logger.info(f'{now_index},{result}')
- result_jsons.append(result)
- else:
- logger.warn(f'{fullName} 无数据,休眠0.01秒')
- logger.info(f"{query_time},获取数据数量:{len(result_jsons)},耗时:{datetime.now() - begin_time}")
- if result_jsons:
- df = complete_minute_data_fast(pd.DataFrame(result_jsons), query_time)
- file_name = f'/data/wind_files/headquarter/cms/tmp/{query_time.replace('-','')}/{wind_turbine_type}/{wind_farm_id}/{wind_turbine_id}.parquet'
- os.makedirs(os.path.dirname(file_name), exist_ok=True)
- df.to_parquet(file_name, index=False)
- logger.info(f'{query_time},{fullName}保存成功,耗时:{datetime.now() - begin_time}')
- # ===================== 统一调用示例 =====================
- if __name__ == "__main__":
- app_begin_time = datetime.now()
- try:
- # 只初始化一次客户端,Token全局复用
- client = CmsApiClient()
- import json, multiprocessing, sys
- # file_name = 'wind_data_exists.json'
- file_suffix = sys.argv[1] if len(sys.argv) > 1 else '1'
- file_name = f'wind_data_with_scada_{file_suffix}.json'
- print(file_name)
- with open(file_name, 'r', encoding='utf-8') as f:
- datas = json.load(f)
- bool_has_result = 0
- time_datas = get_prev_days(prev_days=1)
- total = len(datas)
- for query_time in time_datas:
- with multiprocessing.Pool(6, maxtasksperchild=5) as pool:
- pool.starmap(get_day_data, [(v, query_time) for v in datas])
- # total = 1
- # get_day_data(get_prev_day(prev_days=1))
- logger.info(f"执行结束,总耗时:{datetime.now() - app_begin_time}")
- from combine_cms_data import main
- main()
- except Exception as e:
- logger.error(f"执行出错:{str(e)}")
- raise
|