如何在Python的PIconnect中实现带过滤的summaries统计
解决方案:PI Historian过滤后计算7天时间加权平均
针对你遇到的PIconnect summaries()不支持过滤参数,且全量拉取数据会导致服务器崩溃的问题,提供以下三种可行方案:
方案1:PI AF预创建过滤标签(推荐)
在PI Asset Framework(AF)中创建分析标签,提前在服务器端完成过滤逻辑:当机器速度>0时输出目标标签值,否则输出无值(NoValue)。后续直接对这个新标签调用summaries()计算,无需修改原有Python核心代码。
AF分析表达式示例:
If '{MachineSpeedTag}' > 0 Then '{TargetTag}' Else NoValue()
方案2:分批次拉取单周期数据并本地计算
利用recorded_values()支持过滤的特性,按7天为周期分批拉取数据,在本地计算时间加权平均,避免一次性拉取全年数据导致服务器压力过大。
修改后的Python代码:
import pandas as pd import numpy as np from datetime import datetime, timedelta import PIconnect as PI from OSIsoft.AF.PI import PIException import pytz from datetime import date PI.PIConfig.DEFAULT_TIMEZONE = ' '.join(pytz.country_timezones['za']) bed = PI.PIServer(server="xxxx") # 标签定义保持不变 name = [prod_n, ft1002_n, ft2005_n, ft9000_n, ft8005_n, ft0017_n, ft0035_n, ft5004_n, ft5001_n, ft5002_n, ft8b02_n, ft5003_n, at4001_n, at4005_n, kol_n, dcm_n, hcl_n] = [ bed.search('*u3*tot*prod*')[0], bed.search('*u3*ft1002*')[0], # varswater [L/min] bed.search('*u3*ft2005*')[0], # uitvloeisel [L/min] bed.search('*u3*ft9000*')[0], # totale stoom [t/h] bed.search('*u3*f*8005*')[0], # stoom na BW- en Mkaste [t/h] bed.search('*cd96*f*0017*')[0], # kondensaat terugvoer [L/min] bed.search('*u3*f*0035*')[0], # demin water na BBW kas [L/min] bed.search('*u3*fc5004*')[0], # demin water na berol verdunning [L/min] bed.search('*u3*fc5001*')[0], # demin water na HW storte 1 [L/min] bed.search('*u3*fc5002*')[0], # demin water na HW storte 2 [L/min] bed.search('*u3*fc8b02*mv*')[0], # demin water na H2SO4 verdunning [L/min] bed.search('*u3*fc5003*mv*')[0], # berolvloei [L/min] bed.search('*u3*at4001*md*')[0], # masjienrigting vog [%] bed.search('*u3*at4005*bd9*')[0], # basisgewig [gsm bd] bed.search('*u3*sht*spot*')[0], # kolletjietelling [kolle/m2] bed.search('*u3*sht*dcm*')[0], # DCM uitgetrektes [%] bed.search('*u3*sht*hcl*')[0] # HCl onoplosbares [dpm] ] filter_etiket = bed.search('*u3*mach*speed*')[0] # masjiensnelheid [m/min] def calculate_time_weighted_avg(series): # 计算时间加权平均:(值 * 时间差总和) / 总时间跨度 time_diff = series.index.to_series().diff().dt.total_seconds().fillna(0) weighted_sum = (series.values * time_diff).sum() total_time = time_diff.sum() return weighted_sum / total_time if total_time > 0 else np.nan def dataverwerking(name, filter_etiket, beg, eind): data_list = [] current_start = beg # 遍历所有7天周期 while current_start < eind: current_end = min(current_start + timedelta(days=7), eind) cycle_results = [] for tag in name: try: # 拉取当前周期内过滤后的记录值 df = pd.DataFrame(tag.recorded_values( start_time=current_start, end_time=current_end, filter_expression=f"{filter_etiket} > 0" )) # 计算当前周期的时间加权平均 twa = calculate_time_weighted_avg(df.iloc[:,0]) cycle_results.append(twa) except PIException: cycle_results.append(np.nan) data_list.append((current_start, cycle_results)) current_start = current_end # 整理为与原代码结构一致的DataFrame列表 result_dfs = [] for idx, tag in enumerate(name): tag_data = pd.DataFrame( [(dt, vals[idx]) for dt, vals in data_list], columns=['Timestamp', tag.name] ) tag_data.set_index('Timestamp', inplace=True) result_dfs.append(tag_data) return result_dfs beg = tyd_za("2024.10.01") # FJ 2024 begindatum eind = tyd_za(date.today()) # vandag # 解包数据,保持原代码变量命名 prod, ft1002, ft2005, ft9000, ft8005, ft0017, ft0035, ft5004, ft5001, ft5002, ft8b02, ft5003, at4001, at4005, kol, dcm, hcl = dataverwerking(name, filter_etiket, beg, eind)
方案3:直接调用PI Web API
PIconnect的summaries()未封装过滤参数,但PI Web API的/streams/{WebId}/summaries端点支持filter参数,可直接构造HTTP请求实现服务器端过滤后计算。
示例代码(使用requests库):
import requests import pandas as pd from datetime import datetime # PI Web API配置 base_url = "https://your-pi-web-api-server/piwebapi" pi_server_name = "xxxx" auth = ("your_username", "your_password") def get_tag_web_id(tag_path): url = f"{base_url}/points?path=\\\\{pi_server_name}\\{tag_path}" response = requests.get(url, auth=auth) return response.json()["WebId"] def get_filtered_summaries(tag_web_id, filter_tag_path, start_time, end_time, interval="7d"): summary_url = f"{base_url}/streams/{tag_web_id}/summaries" params = { "startTime": start_time.strftime("%Y-%m-%dT%H:%M:%SZ"), "endTime": end_time.strftime("%Y-%m-%dT%H:%M:%SZ"), "interval": interval, "summaryType": "TimeWeightedAverage", "filter": f"\\\\{pi_server_name}\\{filter_tag_path} > 0", "calculationBasis": "TimeWeighted" } response = requests.get(summary_url, params=params, auth=auth) data = response.json() # 解析结果为DataFrame df = pd.DataFrame([ { "Timestamp": entry["Timestamp"], "TimeWeightedAverage": entry["Value"]["Value"] } for entry in data["Items"] ]) df["Timestamp"] = pd.to_datetime(df["Timestamp"]) df.set_index("Timestamp", inplace=True) return df # 使用示例 filter_tag_path = "YOUR_MACHINE_SPEED_TAG_PATH" target_tag_path = "YOUR_TARGET_TAG_PATH" beg = tyd_za("2024.10.01") eind = tyd_za(date.today()) tag_web_id = get_tag_web_id(target_tag_path) df = get_filtered_summaries(tag_web_id, filter_tag_path, beg, eind)
内容的提问来源于stack exchange,提问作者Philip de Bruin
相关产品推荐
相关产品推荐

