You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 10:27:03