| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515 |
- #!/usr/bin/env python3
- # -*- coding: utf-8 -*-
- r"""变桨面生成端 —— 补上"包内没有生成端"的 `pitch/**` 两件 (用户令 2026-09-19)。
- ## 为什么要写它
- 用户令:「**所有的计算 均要形成 观澜的源代码**,确保 观澜 在别的电脑安装后,系统运行正常。」
- `pitch/pitch_daily.parquet` 与 `pitch/pitch_zero_monthly.parquet` 是**随包快照**(族表里 kind=shipped、
- gen=None),于是另一台机器上装完、放数据、重算之后,变桨面永远只能显示"缺件 ⇒ 不可判"。
- 本脚本把它**从 `data/raw` 算出来**(与其它 20 多个构建器同一套纪律:生成端进包 + 进重算链 + 产物自登记)。
- ## 产物与消费者 (字段契约见 docs/源代码化清单_v0.1.md §2)
- pitch/pitch_daily.parquet date, turbine, n, hydlevel_timeon, hydfilt_timeon, hyd_band,
- hyd_over_relief, hyd_under_pump, hublub_timeon, pitchpum_timeon,
- pistA, pistB, pistC, n_op, hyd_mean
- → src/windscada/subsys/pitch.py::registry 四维度判级
- pitch/pitch_zero_monthly.parquet month, turbine, zero_dev, n_rows + 两口径分列/佐证列
- → 同上 D3「执行与位置反馈」(零位轴)
- pitch/pitch_face_manifest.json 口径/阈值/证据/计数 (人可核对;也是"不静默造数"的凭据)
- ## 口径 (用户令 2026-09-19:「零位取停机段中位、满发段中位两者**分列**出」)
- **来源分工**(实测:两组通道互不重叠):
- scada_10min/<WTGnn>.csv 开关量 timeon / 压力分位 / 运行状态 → 日粒度一行 (n×600 = 当日应测秒)
- scada_1min/<台号>.csv 三叶桨距角 / 有功 / 转速 / 液压压力 → 零位两口径 + 越线分钟计数
- `pitch_daily` 各列(**每列都是"当日"聚合**,不是窗内均值):
- n 当日 10min 记录数 (满日 144;×600 = 应测秒 —— 消费者按 `n*600 - timeon` 算报警秒)
- n_op 当日**运行段数** (判据 flg_wtc_ScInOper_endvalue==1;缺列则退 tur_wtc_GenRpm_mean>0.5)
- ★消费者的不可判门是 `窗内 nop < 144`(≈不足 1 天等效) ⇒ n_op 必须与 n 同粒度
- hyd_mean 当日 prs_wtc_HydPress_mean 均值 (bar)
- hyd_band 当日**运行段**段内 (HydPress_max − HydPress_min) 的**中位** = 蓄能器锯齿带宽
- hyd_over_relief 当日运行**分钟**中 系统压力 ≥ 卸荷阈(240bar) 的分钟数
- hyd_under_pump 当日运行**分钟**中 系统压力 ≤ 启泵阈(195bar) 的分钟数
- hublub_timeon 当日 dot_wtc_HubLubPu_timeon 求和 (s/日) ← 轮毂润滑泵
- pitchpum_timeon 当日 dot_wtc_PitchPum_timeon 求和 (s/日) ← 变桨泵
- hydlevel_timeon 当日 din_wtc_HydLevel_timeon 求和 (s/日) ← 油位开关(1=正常 ⇒ 报警秒=应测−timeon)
- hydfilt_timeon 当日 din_wtc_HydrFilt_timeon 求和 (s/日) ← 滤网开关
- pistA/B/C 当日 din_wtc_LubPist{A,B,C}_counts 求和 (次/日) ← 三叶柱塞动作次数(不对称轴)
- **阈值出处(不是拍的,是现场数据上读出来的)**:该场系统压力锯齿的下沿中位 **194 bar**、
- 上沿中位 **245~246 bar**(10min 段内 max/min 的 fleet 中位);1min 运行分钟压力 **P5 = 195 / P95 = 240**。
- 取整定为 **卸荷阈 240 bar / 启泵阈 195 bar**:① 与"压力越过正常锯齿上下沿"的物理含义一致;
- ② 计数量级与旧随包件同(旧件 over 36~130 次/日、under 0~40 次/日)。
- ★ 与旧件的差异(如实标注,不假装复刻):旧件 `n_op` 口径未知(实测旧件 WTG01 2025-01-01
- n_op=48,本器同日在位运行段 61 段),故**旧件只作量级对拍、不作逐值验收**。
- `pitch_zero_monthly` 两口径(**1min** 数据,月 × 台):
- 停机段中位 该月「顺桨停机」分钟(转速 < 0.5 rpm ∧ |有功| < 50 kW ∧ 三叶桨距角中位 ∈ 顺桨带 [85,92]°)
- 三叶桨距角中位 − 同月**全场中位基准**
- 满发段中位 该月「满发」分钟(有功 ≥ 95% 额定 = 3800 kW)三叶桨距角中位 − 同月**全场中位基准**
- 运行段分档 该月「运行」分钟(转速 ≥ 0.5 rpm ∧ 有功 100~3800 kW)按有功 500kW 分档, 逐档比同月全场
- 同档中位, 再取该台跨档中位 (可比档 ≥4 档、单档 ≥30 分钟) —— **判据轴**
- 物理含义与为什么是**三**个口径(第三口径是 2026-09-19 反推旧件后补的, 有实测依据):
- ① 停机顺桨段 = 叶片压在**机械顺桨止挡**上的读数差。止挡是机械硬基准, fleet 顺桨位集中在
- 88.59~88.69°(±0.7°) ⇒ 这个口径量的是"止挡/编码器在该位的差", **量不出运行区零位差**:
- 实测它把 19#(现场 −0.87° 专项确诊) 读成 **+0.11°**。
- ② 满发段 = 守住额定功率所加的桨角差; 受风况/空气密度影响大(台间离散大), 只作**并列参考**。
- ③ 运行段分档 = 同工况下桨距角读数差, 即行业口径的"零位偏差"; 分档是为了不让风资源混进来
- (风大的台长期在高桨距档 ⇒ 中位天然偏高)。实测它**把 19# 单独拎出来**(逐月持续为负, 其余台 ≈0),
- 与现场专项确诊断性一致 ⇒ `zero_dev_run` 作判据轴, 另两口径并列可见。
- `zero_dev` 保留=停机段(用户令选定的列), 消费者 `pitch.py` 按"两口径取更不利"用。
- 旧随包件的 `zero_dev` 是"绝对零位偏差"(实测只有 19# 离群 −0.896、其余台 ≈0), 与本器"相对同月全场中位"
- 不同源 ⇒ 零位一族**不做逐值对拍**, 只作定性和量级对照(见 manifest 证据「零位三口径」)。
- 有效门(不静默造数):该台该月顺桨停机分钟 < 50 ⇒ 不出值;同月全场有效台 < 5 ⇒ 该月无基准,
- 整月不出值;运行段可比档 < 4 ⇒ `zero_dev_run` 留空(不拿 fleet 值顶替)。三叶极差中位
- (`spread_stop/full`) 一并存下(叶间不平衡证据,消费者暂未用)。
- ## 用法
- python scripts/pitch_face_build.py # 算并写盘 (38 台全量, 约 10 分钟)
- python scripts/pitch_face_build.py --turbines WTG01,WTG19 # 只算几台 (调试)
- python scripts/pitch_face_build.py --dry-run # 只报会写什么、源件在哪、能不能算
- python scripts/pitch_face_build.py --check-map # 复核 1min 台号 → WTG 台号 的映射 (对拍功率)
- """
- from __future__ import annotations
- import argparse
- import json
- import pathlib
- import re
- import sys
- import time
- ROOT = pathlib.Path(__file__).resolve().parents[1]
- sys.path.insert(0, str(ROOT))
- # ★GBK 控制台/日志重定向不再因 '²' '⇒' 这类字符抛 UnicodeEncodeError (2026-09-19 远端实逮:
- # `baseline_38_build.py` 打印 'm/s²' 时在 cp936 下 rc=1, 整条恢复链少一件产物)。
- # 只把**不可编码字符替换掉**, 不改流编码 —— 中文照旧可读。
- try:
- sys.stdout.reconfigure(errors='replace')
- except Exception:
- pass
- import numpy as np # noqa: E402
- import pandas as pd # noqa: E402
- from src import paths as P # noqa: E402
- from src.derived_manifest import record as dm_record # noqa: E402
- from src.windscada.config import farm # noqa: E402
- from src.windscada import scada_source as SS # noqa: E402
- # ── 阈值与门 (出处见模块 docstring「阈值出处」) ──────────────────────────────────
- HYDR_UNLOAD_BAR = 240.0 # 越卸荷阈 (bar): 1min 运行分钟压力 P95 = 240
- HYDR_PUMP_BAR = 195.0 # 跌启泵阈 (bar): 1min 运行分钟压力 P5 = 195
- STOP_RPM = 0.5 # 停机判据: 转速 < 0.5 rpm
- STOP_KW = 50.0 # 停机判据: |有功| < 50 kW
- FULL_FRAC = 0.95 # 满发判据: 有功 ≥ 95% 额定 (4000 kW → 3800 kW)
- # 顺桨停机带 (°): 「停机」分钟里三叶桨距角**是双峰的** —— 实测 2025-01 停机分钟的 5/25/50/75/95 分位 =
- # 28.7 / 31.2 / **31.79** / 88.59 / 88.63: 一半停在顺桨位(≈88.6°), 一半停在 ~31.8°(未顺桨/维护位)。
- # 零位基准只能取**压住机械顺桨止挡**的那一支, 故按带滤: fleet 顺桨位实测 88.59~88.69°(±0.7°),
- # 带取 85~92°(±3.5° 容差) —— 既滤掉 31.8° 那一支, 又给真实零位误差留 ±3.4° 的可见余量。
- FEATHER_BAND = (85.0, 92.0)
- MIN_STOP_MIN = 50 # 台月顺桨停机段最少分钟数
- MIN_BASE_TURBINES = 5 # 同月基准至少要有几台有效台
- # ── 运行段(细桨距)零位口径 —— **工况分档归一** (第三个口径, 2026-09-19 反推旧件后补) ──
- # 零位偏差的行业含义是"**同一工况**下桨距角读数差"; 不分档会被风资源混入 (风大的台长期在高桨距档
- # ⇒ 中位天然偏高)。故按有功 500kW 分档, 逐档与全场同月中位比, 再取该台跨档中位 = `zero_dev_run`。
- # 为什么必须有它 (实测): 停机顺桨段口径把 19# 读成 +0.11° (顺桨止挡是**机械基准**, 与运行区零位
- # 不是同一个量; 且 fleet 顺桨位 88.59~88.69° 极紧, 差异被止挡抹平), 满发段口径 19# = +0.61° 且台间
- # 离散大; 而**运行段口径把 19# 单独拎出来** (逐月持续为负, 其余台 ≈ 0) —— 与现场
- # "19# 零位 −0.87° 专项确诊" 定性一致 ⇒ 判据轴不能只靠停机段那一个口径。
- RUN_MIN_KW = 100.0 # 运行段下界 (kW): 低于此视为停机/待机
- RUN_BIN_KW = 500.0 # 工况分档宽 (kW)
- MIN_RUN_BIN_N = 30 # 单档单月最少分钟数
- MIN_RUN_BINS = 4 # 台月至少要有几档可比
- # ── 列清单 (源件列名 → 本器内部名) ─────────────────────────────────────────────
- C10 = dict(
- state='flg_wtc_ScInOper_endvalue', rpm='tur_wtc_GenRpm_mean',
- hp_min='prs_wtc_HydPress_min', hp_max='prs_wtc_HydPress_max', hp_mean='prs_wtc_HydPress_mean',
- lvl='din_wtc_HydLevel_timeon', filt='din_wtc_HydrFilt_timeon',
- hublub='dot_wtc_HubLubPu_timeon', pitchpum='dot_wtc_PitchPum_timeon',
- pA='din_wtc_LubPistA_counts', pB='din_wtc_LubPistB_counts', pC='din_wtc_LubPistC_counts')
- TS10 = 'TimeStamp'
- C1 = dict(ts='time_stamp', pw='active_power', rpm='rotor_speed', pr='液压站预充压力',
- b1='pitch_angle_blade_1', b2='pitch_angle_blade_2', b3='pitch_angle_blade_3')
- DAILY_COLS = ['date', 'turbine', 'n', 'n_op', 'hyd_mean', 'hyd_band', 'hyd_over_relief', 'hyd_under_pump',
- 'hublub_timeon', 'pitchpum_timeon', 'hydlevel_timeon', 'hydfilt_timeon',
- 'pistA', 'pistB', 'pistC']
- ZERO_COLS = ['month', 'turbine', 'zero_dev', 'n_rows', 'zero_dev_stop', 'zero_dev_full',
- 'zero_dev_run', 'run_bins', 'run_med_bin',
- 'stop_med', 'full_med', 'n_stop', 'n_stop_all', 'n_full', 'spread_stop', 'spread_full',
- 'base_stop', 'base_full']
- def turbine_of_1min(stem: str) -> str | None:
- """1min 导出件的**内部台号** → 场配置台号: `01E.csv` → `WTG01`。
- 现场 1min 导出用内部代号 (`01E.csv … 38B.csv`,`scripts/place_raw_data.py`);数字前缀 = 台号。
- 实测佐证 (`--check-map`): 4 组对拍 (01E→WTG01 / 02E→WTG02 / 19C→WTG19 / 38B→WTG38) 的 10min 功率
- 与 1min 功率**相关系数 0.982~0.986、均值差 < 1 kW**;跨台对拍则为不相关 ⇒ 映射按数字前缀成立。
- """
- m = re.match(r'^(\d+)', stem)
- if not m:
- return None
- return f'WTG{int(m.group(1)):02d}'
- def _paths(cfg):
- csv10 = pathlib.Path(str(cfg['src_10min']))
- csv1, _ = SS.dirs_of('1min', cfg)
- return csv10, csv1
- def read_10min(wtg: str, cfg) -> pd.DataFrame:
- """10min 取一台 (CSV 优先 / MDB 回落, 见 src/windscada/scada_source.py)。"""
- cols = [TS10] + list(C10.values())
- d = SS.load('10min', wtg, columns=cols, cfg=cfg)
- have = {v: k for k, v in C10.items()}
- missing = [v for v in C10.values() if v not in d.columns]
- if TS10 not in d.columns and len(d.columns):
- d = d.rename(columns={d.columns[0]: TS10})
- if missing:
- d.attrs['pitch_missing'] = missing
- return d.rename(columns=have)
- def read_1min(code: str, cfg) -> pd.DataFrame:
- """1min 取一台 (CSV 优先 / MDB 回落)。列名同 `C1` 的值;缺列由调用方判缺。"""
- cols = list(C1.values())
- d = SS.load('1min', code, columns=cols, cfg=cfg)
- if C1['ts'] not in d.columns and len(d.columns):
- d = d.rename(columns={d.columns[0]: C1['ts']})
- return d
- def _num(s) -> pd.Series:
- return pd.to_numeric(s, errors='coerce')
- def daily_of(wtg: str, d10: pd.DataFrame) -> pd.DataFrame:
- """10min 表 → 日粒度一行 (列见 DAILY_COLS)。"""
- d = d10.copy()
- d['ts'] = pd.to_datetime(d[TS10], errors='coerce')
- d['date'] = d['ts'].dt.normalize()
- for k in ('hp_min', 'hp_max', 'hp_mean', 'lvl', 'filt', 'hublub', 'pitchpum', 'pA', 'pB', 'pC', 'rpm'):
- if k in d:
- d[k] = _num(d[k])
- if 'state' in d:
- run = _num(d['state']) == 1
- else: # 兜底: 无运行状态列时按发电机转速判运行
- run = d['rpm'] > 0.5
- d['_run'] = run.fillna(False)
- d['_band'] = (d['hp_max'] - d['hp_min']).where(d['_run']) if {'hp_max', 'hp_min'} <= set(d.columns) else np.nan
- g = d.groupby('date')
- spec = dict(n=('ts', 'size'), n_op=('_run', 'sum'), hyd_band=('_band', 'median'))
- for src, dst, how in (('hp_mean', 'hyd_mean', 'mean'), ('lvl', 'hydlevel_timeon', 'sum'),
- ('filt', 'hydfilt_timeon', 'sum'), ('hublub', 'hublub_timeon', 'sum'),
- ('pitchpum', 'pitchpum_timeon', 'sum'), ('pA', 'pistA', 'sum'),
- ('pB', 'pistB', 'sum'), ('pC', 'pistC', 'sum')):
- if src in d.columns: # 缺列 = NaN (不静默造数; 缺列记进 manifest)
- spec[dst] = (src, (lambda s: s.sum(min_count=1)) if how == 'sum' else how)
- out = g.agg(**spec)
- out.insert(0, 'turbine', wtg)
- return out.reset_index()
- def minute_counts_of(wtg: str, d1: pd.DataFrame, rated_kw: float) -> pd.DataFrame:
- """1min 表 → (日 × 台) 的越卸荷/跌启泵**分钟数** + (月 × 台) 的零位两口径原料。"""
- d = d1.copy()
- d['ts'] = pd.to_datetime(d[C1['ts']], errors='coerce')
- d['date'] = d['ts'].dt.normalize()
- d['month'] = d['ts'].dt.strftime('%Y-%m')
- pw, rpm = _num(d[C1['pw']]), _num(d[C1['rpm']])
- pr = _num(d[C1['pr']]) if C1['pr'] in d else pd.Series(np.nan, index=d.index)
- run = (pw.abs() >= STOP_KW) | (rpm >= STOP_RPM)
- d['_over'] = (run & (pr >= HYDR_UNLOAD_BAR)).astype('int64')
- d['_under'] = (run & (pr <= HYDR_PUMP_BAR)).astype('int64')
- day = d.groupby('date').agg(hyd_over_relief=('_over', 'sum'), hyd_under_pump=('_under', 'sum')).reset_index()
- day.insert(0, 'turbine', wtg)
- # 零位两口径 (逐月)
- blades = [_num(d[C1[b]]) for b in ('b1', 'b2', 'b3') if C1[b] in d]
- pm = sum(blades) / len(blades) if blades else pd.Series(np.nan, index=d.index)
- spread = (pd.concat(blades, axis=1).max(axis=1) - pd.concat(blades, axis=1).min(axis=1)) if len(blades) == 3 \
- else pd.Series(np.nan, index=d.index)
- stop_all = (rpm < STOP_RPM) & (pw.abs() < STOP_KW)
- stop = stop_all & pm.between(*FEATHER_BAND) # 只取压住顺桨止挡的那一支 (见 FEATHER_BAND 注释)
- full = pw >= FULL_FRAC * rated_kw
- run = (rpm >= STOP_RPM) & (pw >= RUN_MIN_KW) & (pw < FULL_FRAC * rated_kw)
- recs = []
- for m, idx in d.groupby('month').groups.items():
- i = pd.Index(idx)
- s, f = stop.loc[i], full.loc[i]
- recs.append(dict(
- month=m, turbine=wtg,
- n_stop=int(s.sum()), n_stop_all=int(stop_all.loc[i].sum()), n_full=int(f.sum()),
- stop_med=float(pm.loc[i][s].median()) if int(s.sum()) else np.nan,
- full_med=float(pm.loc[i][f].median()) if int(f.sum()) else np.nan,
- spread_stop=float(spread.loc[i][s].median()) if int(s.sum()) else np.nan,
- spread_full=float(spread.loc[i][f].median()) if int(f.sum()) else np.nan,
- run_med=float(pm.loc[i][run.loc[i]].median()) if int(run.loc[i].sum()) else np.nan))
- # 运行段按有功分档 (逐月 × 档): 只有足够样本的档才参与 (见 RUN_BIN_KW 注释)
- rb = pd.DataFrame(dict(month=d['month'], bin=(pw // RUN_BIN_KW), pm=pm))[run]
- rb = rb.dropna(subset=['pm'])
- bins = (rb.groupby(['month', 'bin'])['pm'].agg(['median', 'size'])
- .reset_index().rename(columns={'median': 'med', 'size': 'n'}))
- bins = bins[bins['n'] >= MIN_RUN_BIN_N]
- bins.insert(1, 'turbine', wtg)
- return day, pd.DataFrame(recs), bins
- def run_dev(run_bins: pd.DataFrame) -> pd.DataFrame:
- """运行段分档 → 台月 `zero_dev_run` (逐档比全场同月中位, 再取跨档中位)。
- 为什么要分档 (见 RUN_BIN_KW 注释): 不分档时"桨距角中位"混了风资源 (风大的台长期在高桨距档);
- 分档后每档内工况可比, 差值才归因于零位/标定。台月可比档数 < MIN_RUN_BINS ⇒ 不出值。
- """
- if not len(run_bins):
- return pd.DataFrame(columns=['month', 'turbine', 'zero_dev_run', 'run_bins', 'run_med_bin'])
- b = run_bins.copy()
- b['fleet'] = b.groupby(['month', 'bin'])['med'].transform('median')
- b['dev'] = b['med'] - b['fleet']
- g = b.groupby(['month', 'turbine'])['dev'].agg(['median', 'size']).reset_index()
- g = g[g['size'] >= MIN_RUN_BINS].rename(columns={'median': 'zero_dev_run', 'size': 'run_bins'})
- return g
- def zero_monthly(raw: pd.DataFrame, rdev: pd.DataFrame) -> tuple[pd.DataFrame, list]:
- """逐月原料 → 三口径分列 (相对**同月全场中位基准**的偏差)。返回 (表, 被剔记录)。"""
- rejected = []
- if not len(raw):
- return pd.DataFrame(columns=ZERO_COLS), rejected
- df = raw.copy()
- df['stop_ok'] = (df['n_stop'] >= MIN_STOP_MIN) & df['stop_med'].between(*FEATHER_BAND)
- df['full_ok'] = df['n_full'] >= MIN_STOP_MIN
- bad = df[(df['n_stop'] >= MIN_STOP_MIN) & ~df['stop_ok']]
- for _, r in bad.iterrows():
- rejected.append(dict(month=r['month'], turbine=r['turbine'], why='停机段中位不在顺桨带',
- stop_med=round(float(r['stop_med']), 3), n_stop=int(r['n_stop'])))
- short = df[df['n_stop'] < MIN_STOP_MIN]
- n_short = int(len(short))
- if n_short:
- rejected.append(dict(why='台月顺桨停机分钟不足 (不出值)',
- min_minutes=MIN_STOP_MIN, n_turbine_months=n_short,
- example=[dict(month=r['month'], turbine=r['turbine'], n_stop=int(r['n_stop']),
- n_stop_all=int(r['n_stop_all'])) for _, r in short.head(5).iterrows()]))
- base_stop = df[df['stop_ok']].groupby('month')['stop_med'].agg(['median', 'size'])
- base_full = df[df['full_ok']].groupby('month')['full_med'].agg(['median', 'size'])
- ok_months = set(base_stop[base_stop['size'] >= MIN_BASE_TURBINES].index) & \
- set(base_full[base_full['size'] >= MIN_BASE_TURBINES].index)
- for m in sorted(set(df['month']) - ok_months):
- rejected.append(dict(month=m, why='同月有效台不足或无基准', min_base=MIN_BASE_TURBINES))
- df = df[df['month'].isin(ok_months) & df['stop_ok'] & df['full_ok']].copy()
- if not len(df):
- return pd.DataFrame(columns=ZERO_COLS), rejected
- df['base_stop'] = df['month'].map(base_stop['median'])
- df['base_full'] = df['month'].map(base_full['median'])
- df['zero_dev_stop'] = (df['stop_med'] - df['base_stop']).round(4)
- df['zero_dev_full'] = (df['full_med'] - df['base_full']).round(4)
- df['zero_dev'] = df['zero_dev_stop'] # 消费者读这一列 (口径=停机段, 见 docstring)
- df['n_rows'] = df['n_stop'] # 旧件同名列: 样本数
- if len(rdev): # 第三口径 (运行段分档): 可比档不足 ⇒ NaN (不造数)
- df = df.merge(rdev, on=['month', 'turbine'], how='left')
- df['zero_dev_run'] = df['zero_dev_run'].round(4)
- for c in ('zero_dev_run', 'run_bins', 'run_med_bin'):
- if c not in df.columns:
- df[c] = np.nan
- df = df[ZERO_COLS].sort_values(['turbine', 'month']).reset_index(drop=True)
- return df, rejected
- def check_map(cfg, n: int = 4) -> int:
- """复核 1min 台号 → WTG 台号 的映射 (功率对拍)。返回 rc。"""
- csv10, csv1 = _paths(cfg)
- codes = sorted(f.stem for f in csv1.glob('*.csv'))
- picks = [c for c in codes if turbine_of_1min(c) in ('WTG01', 'WTG02', 'WTG19', 'WTG38')][:n]
- ok = 0
- print(f'映射复核 (1min 台号 → 场台号): 抽 {len(picks)} 组, 判据 = 10min 功率与 1min 功率相关系数 > 0.9')
- for code in picks:
- wtg = turbine_of_1min(code)
- a = read_1min(code, cfg)
- a['ts'] = pd.to_datetime(a[C1['ts']], errors='coerce').dt.floor('10min')
- g = _num(a[C1['pw']]).groupby(a['ts']).mean().rename('p1')
- b = SS.load('10min', wtg, columns=[TS10, 'grd_wtc_ActPower_mean'], cfg=cfg)
- b['ts'] = pd.to_datetime(b[TS10], errors='coerce')
- m = pd.merge(g.reset_index(), b.rename(columns={'grd_wtc_ActPower_mean': 'p10'}), on='ts').dropna()
- r = float(np.corrcoef(m['p1'], _num(m['p10']))[0, 1]) if len(m) > 100 else float('nan')
- good = r > 0.9
- ok += good
- print(f' [{"OK" if good else "X "}] {code}.csv → {wtg}: 重叠 {len(m)} 段 r={r:.4f} '
- f'1min 均值 {m["p1"].mean():.1f} kW / 10min 均值 {_num(m["p10"]).mean():.1f} kW')
- return 0 if ok == len(picks) and picks else 1
- def main() -> int:
- ap = argparse.ArgumentParser(description='变桨面生成端 (pitch_daily + pitch_zero_monthly)')
- ap.add_argument('--farm', default='rudong')
- ap.add_argument('--turbines', default=None, help='逗号分隔, 只算这几台 (调试)')
- ap.add_argument('--dry-run', action='store_true', help='只报源件/口径, 不写盘')
- ap.add_argument('--check-map', action='store_true', help='复核 1min 台号 → WTG 台号 映射')
- args = ap.parse_args()
- cfg = farm(args.farm)
- if args.check_map:
- return check_map(cfg)
- rated = float(cfg.get('rated_kw') or 4000)
- csv10, csv1 = _paths(cfg)
- store, pitch_dir = P.store(args.farm), P.pitch(args.farm)
- turb_cfg = [f'WTG{i:02d}' for i in range(1, int(cfg.get('n_turbines') or 38) + 1)]
- want = [t.strip() for t in args.turbines.split(',')] if args.turbines else turb_cfg
- print(f'变桨面生成端 · 场={args.farm} 额定={rated:.0f}kW 台数={len(want)}')
- print(f' 10min 源: {csv10} ({len(list(csv10.glob("*.csv")))} 件 CSV, MDB 回落见 scada_source)')
- print(f' 1min 源: {csv1} ({len(list(csv1.glob("*.csv")))} 件 CSV)')
- code_of = {}
- for f in sorted(csv1.glob('*.csv')):
- t = turbine_of_1min(f.stem)
- if t:
- code_of[t] = f.stem
- daily, counts, zeros, zbins, skipped, t0 = [], [], [], [], [], time.time()
- for i, wtg in enumerate(want, 1):
- try:
- d10 = read_10min(wtg, cfg)
- except Exception as e: # 该台源件缺失: 如实记, 不造数
- skipped.append(dict(turbine=wtg, why=f'10min 取数失败: {type(e).__name__}: {e}'[:200]))
- print(f' [{i:2d}/{len(want)}] {wtg}: 跳过 — {skipped[-1]["why"][:80]}', flush=True)
- continue
- if len(d10):
- daily.append(daily_of(wtg, d10))
- for c in (d10.attrs.get('pitch_missing') or []):
- skipped.append(dict(turbine=wtg, why=f'10min 缺列 {c}'))
- code = code_of.get(wtg)
- n1 = 0
- if code:
- try:
- d1 = read_1min(code, cfg)
- day1, z1, b1_ = minute_counts_of(wtg, d1, rated)
- n1 = len(d1)
- counts.append(day1)
- zeros.append(z1)
- if len(b1_):
- zbins.append(b1_)
- except Exception as e:
- skipped.append(dict(turbine=wtg, why=f'1min ({code}) 取数失败: {type(e).__name__}: {e}'[:200]))
- else:
- skipped.append(dict(turbine=wtg, why='1min 无该台导出件 (越线计数/零位三口径按缺)'))
- print(f' [{i:2d}/{len(want)}] {wtg}: 10min {len(d10):6d} 行'
- + (f' · 1min {n1:7d} 行 ({code})' if code else ' · 1min 缺件')
- + f' 已用 {time.time() - t0:.0f}s', flush=True)
- if not daily:
- print('[X] 一件日粒度都没算出来 —— 先放原始件 (scripts/place_raw_data.py) 再跑')
- return 4
- dl = pd.concat(daily, ignore_index=True)
- dl['date'] = pd.to_datetime(dl['date'])
- if counts: # 越卸荷/跌启泵计数来自 1min 侧: 按 (台, 日) 左连接
- cc = pd.concat(counts, ignore_index=True)
- cc['date'] = pd.to_datetime(cc['date'])
- dl = dl.merge(cc, on=['turbine', 'date'], how='left')
- for c in DAILY_COLS:
- if c not in dl.columns:
- dl[c] = np.nan
- dl['n'] = dl['n'].fillna(0).astype('int64')
- dl['n_op'] = dl['n_op'].fillna(0).astype('int64')
- dl = dl[DAILY_COLS].sort_values(['turbine', 'date']).reset_index(drop=True)
- zr = pd.concat(zeros, ignore_index=True) if zeros else pd.DataFrame()
- rb_all = pd.concat(zbins, ignore_index=True) if zbins else pd.DataFrame()
- zm, rejected = zero_monthly(zr, run_dev(rb_all))
- span = (str(dl['date'].min().date()), str(dl['date'].max().date())) if len(dl) else ('', '')
- print(f'\n日粒度 {len(dl)} 行 · 台 {dl["turbine"].nunique()} · 跨度 {span[0]} → {span[1]}')
- print(f'零位月表 {len(zm)} 行 · 台 {zm["turbine"].nunique() if len(zm) else 0} · '
- f'月 {sorted(zm["month"].unique())[:3]}…{sorted(zm["month"].unique())[-2:] if len(zm) else ""}')
- if len(zm):
- idx = zm.set_index(['turbine', 'month'])
- for t in ('WTG19', 'WTG31', 'WTG01'):
- if t not in idx.index.get_level_values(0):
- continue
- r = idx.loc[t]
- print(f' 锚点 {t}: 停机顺桨段 {r["zero_dev_stop"].median():+.3f}° · 满发段 {r["zero_dev_full"].median():+.3f}° · '
- f'运行段分档 {r["zero_dev_run"].median():+.3f}° (可比档中位 {r["run_bins"].median():.0f})')
- for c, tag in (('zero_dev', '停机顺桨段'), ('zero_dev_full', '满发段'), ('zero_dev_run', '运行段分档')):
- s = zm[c].dropna()
- if len(s):
- print(f' 口径 {tag:8s}: |偏差| 中位 {s.abs().median():.3f}° · P95 {s.abs().quantile(.95):.3f}° · '
- f'最大 {s.abs().max():.3f}° · 有值 {len(s)} 台月')
- if skipped:
- print(f' 跳过/缺件 {len(skipped)} 条: ' + '; '.join(f'{s["turbine"]} {s["why"][:40]}' for s in skipped[:3])
- + (' …' if len(skipped) > 3 else ''))
- if rejected:
- print(f' 未出值(口径门) {len(rejected)} 条: '
- + '; '.join(f'{r.get("month", "—")} {r.get("turbine", "")} {r["why"]}' for r in rejected[:3])
- + (' …' if len(rejected) > 3 else ''))
- if args.dry_run:
- print('\n[dry-run] 不写盘。')
- return 0
- pitch_dir.mkdir(parents=True, exist_ok=True)
- f_daily, f_zero = pitch_dir / 'pitch_daily.parquet', pitch_dir / 'pitch_zero_monthly.parquet'
- dl.to_parquet(f_daily, index=False)
- zm.to_parquet(f_zero, index=False)
- man = dict(
- meta=dict(built=time.strftime('%Y-%m-%d %H:%M:%S'), farm=args.farm, by='scripts/pitch_face_build.py',
- version='pitch-face-v1', turbines=int(dl['turbine'].nunique()),
- span=dict(from_=span[0], to=span[1]), rows=dict(daily=int(len(dl)), zero_monthly=int(len(zm)))),
- 口径=dict(
- n='当日 10min 记录数 (×600 = 应测秒)', n_op='当日运行段数 (flg_wtc_ScInOper_endvalue==1)',
- hyd_mean='当日 HydPress_mean 均值 (bar)', hyd_band='当日运行段段内 (max−min) 中位 (bar)',
- hyd_over_relief=f'当日运行分钟中 压力 ≥ {HYDR_UNLOAD_BAR:.0f}bar 的分钟数',
- hyd_under_pump=f'当日运行分钟中 压力 ≤ {HYDR_PUMP_BAR:.0f}bar 的分钟数',
- pistA_B_C='当日 LubPist{A,B,C}_counts 求和 (三叶柱塞动作次数)',
- zero_dev=f'顺桨停机段中位 − 同月全场中位基准 (停机=转速<{STOP_RPM}rpm ∧ |有功|<{STOP_KW:.0f}kW '
- f'∧ 桨距角中位 ∈ 顺桨带 {FEATHER_BAND[0]:.0f}~{FEATHER_BAND[1]:.0f}°)',
- zero_dev_full=f'满发段中位 − 同月全场中位基准 (满发=有功≥{FULL_FRAC:.0%}额定={FULL_FRAC * rated:.0f}kW)',
- zero_dev_run=f'运行段按有功 {RUN_BIN_KW:.0f}kW 分档, 逐档比同月全场同档中位, 再取跨档中位 '
- f'(判据轴; 运行段=转速≥{STOP_RPM}rpm ∧ 有功 {RUN_MIN_KW:.0f}~{FULL_FRAC * rated:.0f}kW; '
- f'可比档 ≥{MIN_RUN_BINS} 档、单档 ≥{MIN_RUN_BIN_N} 分钟)'),
- 阈值=dict(卸荷阈bar=HYDR_UNLOAD_BAR, 启泵阈bar=HYDR_PUMP_BAR, 停机转速rpm=STOP_RPM, 停机功率kW=STOP_KW,
- 满发门kW=FULL_FRAC * rated, 顺桨带=list(FEATHER_BAND),
- 停机段最少分钟=MIN_STOP_MIN, 同月基准最少台=MIN_BASE_TURBINES,
- 运行段下界kW=RUN_MIN_KW, 工况档宽kW=RUN_BIN_KW, 单档最少分钟=MIN_RUN_BIN_N,
- 最少可比档=MIN_RUN_BINS),
- 证据=dict(
- 压力锯齿='10min 段内 min 中位 194bar / max 中位 245-246bar; 1min 运行分钟 P5=195 / P95=240',
- 顺桨位='顺桨停机分钟三叶桨距角中位 88.6° (fleet, ±0.7°); 停机分钟双峰(另一支 ≈31.8° 未顺桨) '
- '⇒ 零位基准只取顺桨带内那一支',
- 台号映射='1min 内部台号数字前缀 = 场台号 (4 组功率对拍 r=0.982~0.986); 复核: --check-map',
- 与旧随包件='旧件 n_op 口径未知(WTG01 2025-01-01: 旧 48 / 本器 61), 故只作量级对拍不逐值验收; '
- '旧件 zero_dev 只有 19# 是离群(全期 −0.896°, 其余台 ≈0) ⇒ 其口径是"绝对零位偏差", '
- '与本器的"相对同月全场中位"不同源, 故零位一族不做逐值对拍',
- 零位三口径='2026-09-19 反推: 停机顺桨段把 19# 读成 +0.11°(机械止挡抹平差异), 满发段 19# +0.61° '
- '且台间离散大; 运行段分档口径能把 19# 单独拎出(逐月持续为负, 其余台 ≈0) —— '
- '与现场「19# 零位 −0.87° 专项确诊」定性一致, 故判据轴取 running 口径'),
- 来源=dict(_10min='csv 优先 / mdb 回落 (src/windscada/scada_source.py)',
- _1min='csv 优先 / mdb 回落', inputs=[str(csv10), str(csv1)]),
- 跳过=skipped, 未出值=rejected)
- f_man = pitch_dir / 'pitch_face_manifest.json'
- f_man.write_text(json.dumps(man, ensure_ascii=False, indent=1), encoding='utf-8')
- # ★自登记台账的落点与键 (2026-09-19 实逮): 台账正本是 `P.out_root(farm)/_derived_manifest.json`
- # (与其它 8 个生成端一致), 键必须是**相对产物仓根**的路径。此处若写成 P.store()/P.rel() 则会
- # 另建一本 `windscada/_derived_manifest.json` 且键带 `outputs/rudong/` 前缀 —— 台账读不到,
- # 反向呼应审计还会把那本杂账报成"未归类 1 件"(实测 rc=5)。
- out_root = P.out_root(args.farm)
- dm_record(out_root, {(pitch_dir / n).relative_to(out_root).as_posix(): note for n, note in (
- (f_daily.name, 'scripts/pitch_face_build.py (10min 液压/润滑/柱塞 + 1min 越线计数 → 日粒度)'),
- (f_zero.name, 'scripts/pitch_face_build.py (1min 三叶桨距角 → 停机顺桨段/满发段/运行段分档 三口径零位)'),
- (f_man.name, 'scripts/pitch_face_build.py (口径/阈值/证据/跳过清单)'))},
- by='scripts/pitch_face_build.py')
- print(f'\n[OK] 已写 {P.rel(f_daily)} ({f_daily.stat().st_size / 1024:.0f} KB) · '
- f'{P.rel(f_zero)} ({f_zero.stat().st_size / 1024:.0f} KB) · {P.rel(f_man)} · 已自登记')
- return 0
- if __name__ == '__main__':
- sys.exit(main())
|