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