VPS-78: CSG day 费用口径统一 + 断档自愈 (v5) + v3/v5 迁移脚本入库
- compose/pgdb/csg-snapshot-v3.sql: 补提交(此前只存在于工作区,未入 git)
- compose/pgdb/csg-snapshot-v5.sql: 新增
· csg_ladder_cost_raw() 无舍入阶梯费用助手
· csg_daily_snapshot() v5: day 费用改累积边际差分 + round(...,2);
p_month 上移 + 跨月守卫;日期与用量同源 latest_day_kwh
· 一次性归一历史 day cost(UPDATE 3 行,总差 -0.01)
· csg_backfill_missing_days() + job 1011(insert-only 断档自愈)
- runbooks/pgdb-health.md: 新增 check 9 CSG 归档新鲜度
- hosts/pgdb.md: v5 事实 + Known issues 2026-09-22
This commit is contained in:
@@ -0,0 +1,248 @@
|
||||
-- CSG 长期归档修复 (2026-09-21)
|
||||
-- 背景:TimescaleDB 任务 1008 `csg_daily_snapshot()` 自 2026-08-29 创建起 0 成功 / 294 失败。
|
||||
-- 原因:函数签名是零参数 `csg_daily_snapshot()`,而 TimescaleDB 自定义 job 动作
|
||||
-- 按名字调用 `schema.proc(job_id integer, config jsonb)`,解析不到 → 每次报
|
||||
-- `function or procedure "public.csg_daily_snapshot(integer, jsonb)" does not exist`。
|
||||
-- 后果:csg_history 的 day 行停在 2026-08-28、month 行停在 2026-08-01(且 2026-08 的
|
||||
-- month 行是 08-29 当时的「本月至今」302.47 kWh / 180.28 元,不是月终值 331.22 / 198.65)。
|
||||
--
|
||||
-- 本文件的三个动作(幂等,可重复执行):
|
||||
-- 1. `csg_ladder_cost(kwh, month)` —— 阶梯电价助手(广州:夏季 5-10 月 260/600,
|
||||
-- 非夏季 200/400;0.589 / 0.639 / 0.889 元每 kWh)。归档侧的唯一计价来源。
|
||||
-- 2. `csg_daily_snapshot(job_id integer, config jsonb)` v3 —— 兼容 TimescaleDB job 签名;
|
||||
-- 并把 month 行按「数据日期 d 所属月份」归属:
|
||||
-- · d 属于当前月 → usage 取集成本月累计,cost 用阶梯助手重算
|
||||
-- · d 属于上月 → usage 由该月 day 行汇总(不读集成 last_month,避开翻月瞬态:
|
||||
-- 2026-08-31 16:32 实测 last_month_total_usage 曾读到 323.49 这种误值)
|
||||
-- 并加「不降级」保护:day 行汇总结果小于已记录的 month 行时不动。
|
||||
-- 3. 回填 2026-08-29 → 2026-09-19 的 day 行(数据源:08-29..08-31 取 scribe
|
||||
-- states_raw 里 `yesterday_kwh` 观测值减一天;09-01..09-19 取 CSG 集成
|
||||
-- `this_month_by_day`,其中 09-06 = 9.29 是 scribe 漏采的观测)。日费用一律用
|
||||
-- 「当日用电 × 当日所处的边际档位」计算,使整月 day 费用之和 == 阶梯月费用。
|
||||
--
|
||||
-- 回滚(脚本内已自动建快照表 csg_history_bak_20260921):
|
||||
-- BEGIN;
|
||||
-- SELECT delete_job((SELECT job_id FROM timescaledb_information.jobs WHERE proc_name='csg_daily_snapshot'));
|
||||
-- DROP FUNCTION IF EXISTS public.csg_daily_snapshot(integer, jsonb);
|
||||
-- CREATE FUNCTION public.csg_daily_snapshot() RETURNS void LANGUAGE plpgsql AS $f$ ...原函数体... $f$;
|
||||
-- SELECT add_job('public.csg_daily_snapshot'::regproc, INTERVAL '1 day');
|
||||
-- DROP FUNCTION IF EXISTS public.csg_ladder_cost(numeric, date);
|
||||
-- DELETE FROM csg_history; INSERT INTO csg_history SELECT * FROM csg_history_bak_20260921;
|
||||
-- COMMIT;
|
||||
-- (或整体放弃归档:DROP TABLE csg_history + 删任务,见 hosts/pgdb.md 回滚段)
|
||||
|
||||
\set ON_ERROR_STOP on
|
||||
|
||||
BEGIN;
|
||||
|
||||
-- ---------- 0. 回滚快照(只建一次) ----------
|
||||
CREATE TABLE IF NOT EXISTS public.csg_history_bak_20260921 AS TABLE public.csg_history;
|
||||
|
||||
-- ---------- 1. 阶梯电价助手 ----------
|
||||
CREATE OR REPLACE FUNCTION public.csg_ladder_cost(kwh numeric, month_date date)
|
||||
RETURNS numeric
|
||||
LANGUAGE sql
|
||||
IMMUTABLE
|
||||
AS $function$
|
||||
SELECT CASE
|
||||
WHEN kwh IS NULL OR kwh <= 0 THEN 0::numeric
|
||||
WHEN kwh <= (CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 260 ELSE 200 END)
|
||||
THEN round(kwh * 0.589, 2)
|
||||
WHEN kwh <= (CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 600 ELSE 400 END)
|
||||
THEN round((CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 260 ELSE 200 END) * 0.589
|
||||
+ (kwh - (CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 260 ELSE 200 END)) * 0.639, 2)
|
||||
ELSE round((CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 260 ELSE 200 END) * 0.589
|
||||
+ ((CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 600 ELSE 400 END)
|
||||
- (CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 260 ELSE 200 END)) * 0.639
|
||||
+ (kwh - (CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 600 ELSE 400 END)) * 0.889, 2)
|
||||
END;
|
||||
$function$;
|
||||
|
||||
-- ---------- 2. 重建为 TimescaleDB 兼容签名的 v3 + 重挂 job ----------
|
||||
-- 注意:不能直接 DROP 旧函数——TimescaleDB job 1008 持有依赖
|
||||
-- (ERROR: cannot drop public.csg_daily_snapshot because background job 1008 depends on it)。
|
||||
-- 因此:delete_job → 换签名重建函数 → add_job 重挂,并把排程对齐文档口径
|
||||
-- (22:30 Asia/Shanghai = 14:30 UTC;原来实际跑在 22:00 UTC = 06:00 CST,
|
||||
-- 那个时点集成还没发布前一日数据,d 会晚一天)。
|
||||
DO $do$
|
||||
BEGIN
|
||||
IF EXISTS (SELECT 1 FROM timescaledb_information.jobs WHERE job_id = 1008) THEN
|
||||
PERFORM delete_job(1008);
|
||||
END IF;
|
||||
END
|
||||
$do$;
|
||||
|
||||
DROP FUNCTION IF EXISTS public.csg_daily_snapshot();
|
||||
|
||||
CREATE OR REPLACE FUNCTION public.csg_daily_snapshot(
|
||||
job_id integer DEFAULT NULL,
|
||||
config jsonb DEFAULT NULL
|
||||
)
|
||||
RETURNS void
|
||||
LANGUAGE plpgsql
|
||||
AS $function$
|
||||
DECLARE
|
||||
d date;
|
||||
v_usage numeric; v_cost numeric; v_ladder text; v_balance numeric;
|
||||
m_usage numeric; v_tariff numeric;
|
||||
p_month date; p_usage numeric; prev_usage numeric;
|
||||
BEGIN
|
||||
-- 数据日期 = CSG 集成「最新一天」(数据滞后一天);COALESCE 兜底昨天
|
||||
SELECT COALESCE(
|
||||
(SELECT (attributes->>'latest_day_date')::date
|
||||
FROM states_raw s JOIN entities e ON e.id = s.metadata_id
|
||||
WHERE e.entity_id = 'sensor.0800041935246530_latest_day_kwh'
|
||||
ORDER BY s.time DESC LIMIT 1),
|
||||
CURRENT_DATE - 1) INTO d;
|
||||
|
||||
-- 一律取「最新有值行」:轮询先写瞬态 unknown、~4s 后才写有值行
|
||||
SELECT s.value INTO v_usage FROM states_raw s JOIN entities e ON e.id=s.metadata_id
|
||||
WHERE e.entity_id='sensor.0800041935246530_yesterday_kwh' AND s.value IS NOT NULL ORDER BY s.time DESC LIMIT 1;
|
||||
SELECT s.value INTO v_cost FROM states_raw s JOIN entities e ON e.id=s.metadata_id
|
||||
WHERE e.entity_id='sensor.0800041935246530_latest_day_cost' AND s.value IS NOT NULL ORDER BY s.time DESC LIMIT 1;
|
||||
SELECT s.state INTO v_ladder FROM states_raw s JOIN entities e ON e.id=s.metadata_id
|
||||
WHERE e.entity_id='sensor.csg_current_ladder' AND s.state IS NOT NULL AND s.state NOT IN ('unknown','unavailable') ORDER BY s.time DESC LIMIT 1;
|
||||
SELECT s.value INTO v_balance FROM states_raw s JOIN entities e ON e.id=s.metadata_id
|
||||
WHERE e.entity_id='sensor.0800041935246530_balance' AND s.value IS NOT NULL ORDER BY s.time DESC LIMIT 1;
|
||||
SELECT s.value INTO m_usage FROM states_raw s JOIN entities e ON e.id=s.metadata_id
|
||||
WHERE e.entity_id='sensor.0800041935246530_this_month_total_usage' AND s.value IS NOT NULL ORDER BY s.time DESC LIMIT 1;
|
||||
-- 日费用回退:原生 latest_day_cost 从未有值(集成侧该传感器恒为 unknown),
|
||||
-- 用「昨日用电 × 当前档费率」估计;账单口径可用时优先原生
|
||||
SELECT s.value INTO v_tariff FROM states_raw s JOIN entities e ON e.id=s.metadata_id
|
||||
WHERE e.entity_id='sensor.csg_current_ladder_tariff' AND s.value IS NOT NULL ORDER BY s.time DESC LIMIT 1;
|
||||
IF v_cost IS NULL AND v_usage IS NOT NULL AND v_tariff IS NOT NULL THEN
|
||||
v_cost := v_usage * v_tariff;
|
||||
END IF;
|
||||
|
||||
INSERT INTO csg_history (period, kind, usage_kwh, cost, ladder, balance)
|
||||
VALUES (d, 'day', v_usage, v_cost, v_ladder, v_balance)
|
||||
ON CONFLICT (period, kind) DO UPDATE SET
|
||||
usage_kwh = COALESCE(EXCLUDED.usage_kwh, csg_history.usage_kwh),
|
||||
cost = COALESCE(EXCLUDED.cost, csg_history.cost),
|
||||
ladder = COALESCE(EXCLUDED.ladder, csg_history.ladder),
|
||||
balance = COALESCE(EXCLUDED.balance, csg_history.balance),
|
||||
updated_at = now();
|
||||
|
||||
-- month 行按 d 所属月份归属(v3:修掉翻月时用错月份累计的 bug)
|
||||
p_month := date_trunc('month', d)::date;
|
||||
|
||||
IF p_month = date_trunc('month', CURRENT_DATE)::date THEN
|
||||
-- 当月:集成本月累计优先,缺失则退回 day 行汇总;费用一律阶梯重算
|
||||
SELECT COALESCE(m_usage, sum(h.usage_kwh)) INTO p_usage
|
||||
FROM csg_history h
|
||||
WHERE h.kind='day' AND h.period >= p_month AND h.period < (p_month + interval '1 month');
|
||||
ELSE
|
||||
-- 上月/更早:只信该月 day 行(集成翻月时 last_month_* 有瞬态误值)
|
||||
SELECT sum(h.usage_kwh) INTO p_usage
|
||||
FROM csg_history h
|
||||
WHERE h.kind='day' AND h.period >= p_month AND h.period < (p_month + interval '1 month');
|
||||
SELECT h.usage_kwh INTO prev_usage FROM csg_history h
|
||||
WHERE h.kind='month' AND h.period = p_month;
|
||||
IF p_usage IS NULL OR (prev_usage IS NOT NULL AND prev_usage > p_usage) THEN
|
||||
RETURN; -- day 行不完整,保留已有 month 行,不降级
|
||||
END IF;
|
||||
END IF;
|
||||
|
||||
IF p_usage IS NOT NULL AND p_usage > 0 THEN
|
||||
INSERT INTO csg_history (period, kind, usage_kwh, cost)
|
||||
VALUES (p_month, 'month', p_usage, public.csg_ladder_cost(p_usage, p_month))
|
||||
ON CONFLICT (period, kind) DO UPDATE SET
|
||||
usage_kwh = EXCLUDED.usage_kwh,
|
||||
cost = EXCLUDED.cost,
|
||||
updated_at = now();
|
||||
END IF;
|
||||
END $function$;
|
||||
|
||||
-- 重挂每日任务:22:30 Asia/Shanghai(14:30 UTC)
|
||||
DO $do$
|
||||
BEGIN
|
||||
IF NOT EXISTS (SELECT 1 FROM timescaledb_information.jobs WHERE proc_name = 'csg_daily_snapshot') THEN
|
||||
PERFORM add_job(proc => 'public.csg_daily_snapshot'::regproc,
|
||||
schedule_interval => INTERVAL '1 day',
|
||||
initial_start => date_trunc('day', now() AT TIME ZONE 'UTC')
|
||||
+ INTERVAL '14 hours 30 minutes',
|
||||
timezone => 'UTC');
|
||||
END IF;
|
||||
END
|
||||
$do$;
|
||||
|
||||
-- ---------- 3. 回填 day 行 2026-08-29 .. 2026-09-19 + 重算月行 ----------
|
||||
-- 现有 day 行(07-01..08-28 来自首批回填)与新行一起参与「月内累计」窗口,
|
||||
-- 使边际档位计算正确;已有非空 cost 不覆盖。
|
||||
WITH newrows(period, usage_kwh) AS (
|
||||
VALUES
|
||||
(DATE '2026-08-29', 10.99::numeric),
|
||||
(DATE '2026-08-30', 10.03::numeric),
|
||||
(DATE '2026-08-31', 7.73::numeric),
|
||||
(DATE '2026-09-01', 8.50::numeric),
|
||||
(DATE '2026-09-02', 6.68::numeric),
|
||||
(DATE '2026-09-03', 8.74::numeric),
|
||||
(DATE '2026-09-04', 6.76::numeric),
|
||||
(DATE '2026-09-05', 7.76::numeric),
|
||||
(DATE '2026-09-06', 9.29::numeric),
|
||||
(DATE '2026-09-07', 7.24::numeric),
|
||||
(DATE '2026-09-08', 13.00::numeric),
|
||||
(DATE '2026-09-09', 7.98::numeric),
|
||||
(DATE '2026-09-10', 8.16::numeric),
|
||||
(DATE '2026-09-11', 8.75::numeric),
|
||||
(DATE '2026-09-12', 10.48::numeric),
|
||||
(DATE '2026-09-13', 11.11::numeric),
|
||||
(DATE '2026-09-14', 7.22::numeric),
|
||||
(DATE '2026-09-15', 9.94::numeric),
|
||||
(DATE '2026-09-16', 7.19::numeric),
|
||||
(DATE '2026-09-17', 10.32::numeric),
|
||||
(DATE '2026-09-18', 7.77::numeric),
|
||||
(DATE '2026-09-19', 11.79::numeric)
|
||||
),
|
||||
allrows AS (
|
||||
SELECT h.period, h.usage_kwh FROM csg_history h
|
||||
WHERE h.kind='day' AND h.period BETWEEN DATE '2026-07-01' AND DATE '2026-09-19'
|
||||
UNION ALL
|
||||
SELECT n.period, n.usage_kwh FROM newrows n
|
||||
WHERE NOT EXISTS (SELECT 1 FROM csg_history h WHERE h.kind='day' AND h.period = n.period)
|
||||
),
|
||||
calc AS (
|
||||
SELECT period, usage_kwh,
|
||||
CASE WHEN extract(month FROM period) BETWEEN 5 AND 10 THEN 260 ELSE 200 END AS t1,
|
||||
CASE WHEN extract(month FROM period) BETWEEN 5 AND 10 THEN 600 ELSE 400 END AS t2,
|
||||
sum(usage_kwh) OVER (PARTITION BY date_trunc('month', period)
|
||||
ORDER BY period ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS c1
|
||||
FROM allrows
|
||||
),
|
||||
slab AS (
|
||||
SELECT period, usage_kwh, t1, t2, c1, c1 - usage_kwh AS c0 FROM calc
|
||||
),
|
||||
priced AS (
|
||||
SELECT period, usage_kwh,
|
||||
round(
|
||||
(CASE WHEN c1 <= t1 THEN c1*0.589
|
||||
WHEN c1 <= t2 THEN t1*0.589 + (c1-t1)*0.639
|
||||
ELSE t1*0.589 + (t2-t1)*0.639 + (c1-t2)*0.889 END)
|
||||
- (CASE WHEN c0 <= t1 THEN c0*0.589
|
||||
WHEN c0 <= t2 THEN t1*0.589 + (c0-t1)*0.639
|
||||
ELSE t1*0.589 + (t2-t1)*0.639 + (c0-t2)*0.889 END), 2) AS cost,
|
||||
CASE WHEN c1 <= t1 THEN '一档' WHEN c1 <= t2 THEN '二档' ELSE '三档' END AS ladder
|
||||
FROM slab
|
||||
)
|
||||
INSERT INTO csg_history (period, kind, usage_kwh, cost, ladder)
|
||||
SELECT period, 'day', usage_kwh, cost, ladder FROM priced
|
||||
ON CONFLICT (period, kind) DO UPDATE SET
|
||||
usage_kwh = EXCLUDED.usage_kwh,
|
||||
cost = COALESCE(csg_history.cost, EXCLUDED.cost),
|
||||
ladder = COALESCE(csg_history.ladder, EXCLUDED.ladder),
|
||||
updated_at = now();
|
||||
|
||||
-- 月行按 day 行汇总重算(只动 2026-08 / 2026-09;2026-07 及更早的月行来自集成
|
||||
-- by_month,day 行并不覆盖整月,重算会把它算小,故不触碰)
|
||||
INSERT INTO csg_history (period, kind, usage_kwh, cost)
|
||||
SELECT date_trunc('month', h.period)::date, 'month', round(sum(h.usage_kwh),2),
|
||||
public.csg_ladder_cost(round(sum(h.usage_kwh),2), date_trunc('month', h.period)::date)
|
||||
FROM csg_history h
|
||||
WHERE h.kind='day' AND h.period BETWEEN DATE '2026-08-01' AND DATE '2026-09-30'
|
||||
GROUP BY date_trunc('month', h.period)
|
||||
ON CONFLICT (period, kind) DO UPDATE SET
|
||||
usage_kwh = EXCLUDED.usage_kwh,
|
||||
cost = EXCLUDED.cost,
|
||||
updated_at = now();
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,269 @@
|
||||
-- CSG 长期归档 v5 (2026-09-22, VPS-78)
|
||||
--
|
||||
-- 背景:v3 修好了 job(1008 → 1010)并回填,但 csg_daily_snapshot() 写 day 行时用的是
|
||||
-- 「昨日用电 × 当前档费率」**单一费率**且未舍入,而 v3 的回填(07-01..09-19)用的是
|
||||
-- 「当月累积的边际差分」并 round(...,2)。两套口径混在同一列:
|
||||
-- · 实测今天只差**精度**(3 行:2026-08-28 / 09-19 / 09-20,总差 -0.01 元);
|
||||
-- · 但一旦某个**跨档日**由夜间任务写出,就会与历史口径不一致,并破坏
|
||||
-- 「day 费用之和 ≈ 阶梯月费用」这条 v3 建立的验收不变量。
|
||||
-- · 时限:2026-09-22 时 csg_current_ladder_remaining_kwh = 80.6、本月日均 8.97
|
||||
-- → 约 2026-09-29 跨入二档。
|
||||
--
|
||||
-- 本文件的四个动作(幂等,可重复执行):
|
||||
-- 1. csg_ladder_cost_raw(kwh, month) —— **无舍入**的阶梯费用助手。
|
||||
-- v3 回填的算法是 round(精确c1 − 精确c0, 2)。若直接用
|
||||
-- csg_ladder_cost(c1) − csg_ladder_cost(c0),会**二次舍入**(7 月实测 1 天差 0.01),
|
||||
-- 故新增 raw 版,只在最后舍入一次。
|
||||
-- 2. csg_daily_snapshot() v5 —— day cost 改为「当月累积边际差分 + round(...,2)」;
|
||||
-- 并把 p_month 的计算**上移到计费之前**,加**跨月守卫**
|
||||
-- (每月 1 日的 d 落在上月,此时 m_usage 是「新月份」累计,不能作起点,否则为负)。
|
||||
-- 同时把「数据日期」与「当日用量」改为取**同一个实体同一行**
|
||||
-- (sensor.0800041935246530_latest_day_kwh):v3 用 latest_day_kwh 取日期、
|
||||
-- 却用 yesterday_kwh 取数值,两个实体可能错配;且 latest_day_kwh 的归档更完整
|
||||
-- (含 2026-09-06 = 9.29 这条 yesterday_kwh 没有的观测)。
|
||||
-- 3. 一次性归一:把历史 day 行的 cost 按新口径重算(实测仅 3 行变化,总差 -0.01)。
|
||||
-- 4. csg_backfill_missing_days() + job 1011 —— **断档自愈**:只 INSERT 缺失日期,
|
||||
-- ON CONFLICT DO NOTHING,**绝不覆盖既有行**;数据源是 scribe.states_raw 里
|
||||
-- latest_day_kwh 的 (latest_day_date 属性, value) 观测对(该属性 v3 起就有归档)。
|
||||
--
|
||||
-- 回滚(脚本已自动建快照表 csg_history_bak_20260922):
|
||||
-- BEGIN;
|
||||
-- SELECT delete_job((SELECT job_id FROM timescaledb_information.jobs
|
||||
-- WHERE proc_name='csg_backfill_missing_days'));
|
||||
-- DROP FUNCTION IF EXISTS public.csg_backfill_missing_days(integer, jsonb);
|
||||
-- DROP FUNCTION IF EXISTS public.csg_ladder_cost_raw(numeric, date);
|
||||
-- -- v3 的 csg_daily_snapshot 用 compose/pgdb/csg-snapshot-v3.sql 的原文重建即可
|
||||
-- DELETE FROM csg_history; INSERT INTO csg_history SELECT * FROM csg_history_bak_20260922;
|
||||
-- COMMIT;
|
||||
-- (或整体放弃归档:DROP TABLE csg_history + 删 job,见 hosts/pgdb.md 回滚段)
|
||||
|
||||
\set ON_ERROR_STOP on
|
||||
|
||||
BEGIN;
|
||||
|
||||
-- ---------- 0. 回滚快照(只建一次) ----------
|
||||
CREATE TABLE IF NOT EXISTS public.csg_history_bak_20260922 AS TABLE public.csg_history;
|
||||
|
||||
-- ---------- 1. 无舍入阶梯费用助手(差分用) ----------
|
||||
CREATE OR REPLACE FUNCTION public.csg_ladder_cost_raw(kwh numeric, month_date date)
|
||||
RETURNS numeric
|
||||
LANGUAGE sql
|
||||
IMMUTABLE
|
||||
AS $function$
|
||||
SELECT CASE
|
||||
WHEN kwh IS NULL OR kwh <= 0 THEN 0::numeric
|
||||
WHEN kwh <= (CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 260 ELSE 200 END)
|
||||
THEN kwh * 0.589
|
||||
WHEN kwh <= (CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 600 ELSE 400 END)
|
||||
THEN (CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 260 ELSE 200 END) * 0.589
|
||||
+ (kwh - (CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 260 ELSE 200 END)) * 0.639
|
||||
ELSE (CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 260 ELSE 200 END) * 0.589
|
||||
+ ((CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 600 ELSE 400 END)
|
||||
- (CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 260 ELSE 200 END)) * 0.639
|
||||
+ (kwh - (CASE WHEN extract(month FROM month_date) BETWEEN 5 AND 10 THEN 600 ELSE 400 END)) * 0.889
|
||||
END;
|
||||
$function$;
|
||||
|
||||
-- ---------- 2. csg_daily_snapshot() v5 ----------
|
||||
CREATE OR REPLACE FUNCTION public.csg_daily_snapshot(
|
||||
job_id integer DEFAULT NULL,
|
||||
config jsonb DEFAULT NULL
|
||||
)
|
||||
RETURNS void
|
||||
LANGUAGE plpgsql
|
||||
AS $function$
|
||||
DECLARE
|
||||
d date;
|
||||
v_usage numeric; v_cost numeric; v_ladder text; v_balance numeric;
|
||||
m_usage numeric;
|
||||
p_month date; p_usage numeric; prev_usage numeric; c0 numeric; c1 numeric;
|
||||
BEGIN
|
||||
-- 数据日期与当日用量取【同一实体的同一行】:latest_day_kwh
|
||||
-- (v3 的 d 来自它的 latest_day_date,但 v_usage 来自另一个实体 yesterday_kwh)
|
||||
SELECT (s.attributes->>'latest_day_date')::date, s.value
|
||||
INTO d, v_usage
|
||||
FROM states_raw s JOIN entities e ON e.id = s.metadata_id
|
||||
WHERE e.entity_id = 'sensor.0800041935246530_latest_day_kwh'
|
||||
AND s.value IS NOT NULL
|
||||
AND s.attributes->>'latest_day_date' IS NOT NULL
|
||||
ORDER BY s.time DESC LIMIT 1;
|
||||
d := COALESCE(d, CURRENT_DATE - 1);
|
||||
|
||||
-- 一律取「最新有值行」:轮询先写瞬态 unknown、~4s 后才写有值行
|
||||
SELECT s.value INTO v_cost FROM states_raw s JOIN entities e ON e.id=s.metadata_id
|
||||
WHERE e.entity_id='sensor.0800041935246530_latest_day_cost' AND s.value IS NOT NULL ORDER BY s.time DESC LIMIT 1;
|
||||
SELECT s.state INTO v_ladder FROM states_raw s JOIN entities e ON e.id=s.metadata_id
|
||||
WHERE e.entity_id='sensor.csg_current_ladder' AND s.state IS NOT NULL AND s.state NOT IN ('unknown','unavailable') ORDER BY s.time DESC LIMIT 1;
|
||||
SELECT s.value INTO v_balance FROM states_raw s JOIN entities e ON e.id=s.metadata_id
|
||||
WHERE e.entity_id='sensor.0800041935246530_balance' AND s.value IS NOT NULL ORDER BY s.time DESC LIMIT 1;
|
||||
SELECT s.value INTO m_usage FROM states_raw s JOIN entities e ON e.id=s.metadata_id
|
||||
WHERE e.entity_id='sensor.0800041935246530_this_month_total_usage' AND s.value IS NOT NULL ORDER BY s.time DESC LIMIT 1;
|
||||
|
||||
-- 月归属必须在计费之前算:跨月守卫要用它
|
||||
p_month := date_trunc('month', d)::date;
|
||||
|
||||
-- 日费用:与 v3 回填同口径 = 当月累积的边际差分,只在最后舍入一次。
|
||||
-- 原生 latest_day_cost 从未有值(集成侧恒 unknown),故这里是唯一路径。
|
||||
IF v_cost IS NULL AND v_usage IS NOT NULL THEN
|
||||
IF p_month = date_trunc('month', CURRENT_DATE)::date THEN
|
||||
-- 当月日:起点 = 本月累计 − 当日用量
|
||||
IF m_usage IS NOT NULL THEN
|
||||
c1 := m_usage;
|
||||
c0 := m_usage - v_usage;
|
||||
END IF;
|
||||
ELSE
|
||||
-- 上月日(每月 1 日必现,因数据滞后一天):起点 = 该月已归档 day 行之和
|
||||
SELECT COALESCE(sum(h.usage_kwh), 0) INTO c0
|
||||
FROM csg_history h
|
||||
WHERE h.kind='day' AND h.period >= p_month AND h.period < d;
|
||||
c1 := c0 + v_usage;
|
||||
END IF;
|
||||
IF c0 IS NOT NULL AND c1 IS NOT NULL AND c0 >= 0 THEN
|
||||
v_cost := round(public.csg_ladder_cost_raw(c1, p_month)
|
||||
- public.csg_ladder_cost_raw(c0, p_month), 2);
|
||||
END IF;
|
||||
END IF;
|
||||
|
||||
INSERT INTO csg_history (period, kind, usage_kwh, cost, ladder, balance)
|
||||
VALUES (d, 'day', v_usage, v_cost, v_ladder, v_balance)
|
||||
ON CONFLICT (period, kind) DO UPDATE SET
|
||||
usage_kwh = COALESCE(EXCLUDED.usage_kwh, csg_history.usage_kwh),
|
||||
cost = COALESCE(EXCLUDED.cost, csg_history.cost),
|
||||
ladder = COALESCE(EXCLUDED.ladder, csg_history.ladder),
|
||||
balance = COALESCE(EXCLUDED.balance, csg_history.balance),
|
||||
updated_at = now();
|
||||
|
||||
-- month 行按 d 所属月份归属(v3:修掉翻月时用错月份累计的 bug)
|
||||
IF p_month = date_trunc('month', CURRENT_DATE)::date THEN
|
||||
-- 当月:集成本月累计优先,缺失则退回 day 行汇总;费用一律阶梯重算
|
||||
SELECT COALESCE(m_usage, sum(h.usage_kwh)) INTO p_usage
|
||||
FROM csg_history h
|
||||
WHERE h.kind='day' AND h.period >= p_month AND h.period < (p_month + interval '1 month');
|
||||
ELSE
|
||||
-- 上月/更早:只信该月 day 行(集成翻月时 last_month_* 有瞬态误值)
|
||||
SELECT sum(h.usage_kwh) INTO p_usage
|
||||
FROM csg_history h
|
||||
WHERE h.kind='day' AND h.period >= p_month AND h.period < (p_month + interval '1 month');
|
||||
SELECT h.usage_kwh INTO prev_usage FROM csg_history h
|
||||
WHERE h.kind='month' AND h.period = p_month;
|
||||
IF p_usage IS NULL OR (prev_usage IS NOT NULL AND prev_usage > p_usage) THEN
|
||||
RETURN; -- day 行不完整,保留已有 month 行,不降级
|
||||
END IF;
|
||||
END IF;
|
||||
|
||||
IF p_usage IS NOT NULL AND p_usage > 0 THEN
|
||||
INSERT INTO csg_history (period, kind, usage_kwh, cost)
|
||||
VALUES (p_month, 'month', p_usage, public.csg_ladder_cost(p_usage, p_month))
|
||||
ON CONFLICT (period, kind) DO UPDATE SET
|
||||
usage_kwh = EXCLUDED.usage_kwh,
|
||||
cost = EXCLUDED.cost,
|
||||
updated_at = now();
|
||||
END IF;
|
||||
END $function$;
|
||||
|
||||
-- ---------- 3. 一次性归一:历史 day 行 cost 按新口径重算 ----------
|
||||
-- 幂等:只更新与新口径不同的行。实测(写入前 dry-run)仅 3 行变化、总差 -0.01。
|
||||
WITH cum AS (
|
||||
SELECT period, usage_kwh, cost,
|
||||
sum(usage_kwh) OVER (PARTITION BY date_trunc('month', period)
|
||||
ORDER BY period ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS c1
|
||||
FROM csg_history
|
||||
WHERE kind='day' AND usage_kwh IS NOT NULL
|
||||
), priced AS (
|
||||
SELECT period,
|
||||
round(public.csg_ladder_cost_raw(c1, period)
|
||||
- public.csg_ladder_cost_raw(c1 - usage_kwh, period), 2) AS new_cost
|
||||
FROM cum
|
||||
)
|
||||
UPDATE csg_history h
|
||||
SET cost = p.new_cost, updated_at = now()
|
||||
FROM priced p
|
||||
WHERE h.kind='day' AND h.period = p.period
|
||||
AND h.cost IS DISTINCT FROM p.new_cost;
|
||||
|
||||
-- ---------- 4. 断档自愈(只补缺失,不覆盖既有) ----------
|
||||
CREATE OR REPLACE FUNCTION public.csg_backfill_missing_days(
|
||||
job_id integer DEFAULT NULL,
|
||||
config jsonb DEFAULT NULL
|
||||
)
|
||||
RETURNS void
|
||||
LANGUAGE plpgsql
|
||||
AS $function$
|
||||
DECLARE
|
||||
lookback int := COALESCE((config->>'lookback_days')::int, 45);
|
||||
dd date;
|
||||
v_usage numeric;
|
||||
c0 numeric; c1 numeric;
|
||||
v_ladder text;
|
||||
mm date;
|
||||
BEGIN
|
||||
-- 4a. 补缺失的 day 行(只 INSERT,绝不覆盖)
|
||||
FOR dd IN
|
||||
SELECT g::date
|
||||
FROM generate_series(CURRENT_DATE - lookback, CURRENT_DATE - 1, INTERVAL '1 day') g
|
||||
WHERE NOT EXISTS (SELECT 1 FROM csg_history h WHERE h.kind='day' AND h.period = g::date)
|
||||
ORDER BY 1
|
||||
LOOP
|
||||
-- 该数据日的观测:同一实体同时给出日期与数值
|
||||
SELECT s.value INTO v_usage
|
||||
FROM states_raw s JOIN entities e ON e.id = s.metadata_id
|
||||
WHERE e.entity_id = 'sensor.0800041935246530_latest_day_kwh'
|
||||
AND s.value IS NOT NULL
|
||||
AND s.attributes->>'latest_day_date' = dd::text
|
||||
ORDER BY s.time DESC LIMIT 1;
|
||||
CONTINUE WHEN v_usage IS NULL; -- 上游没有这天的数据,不写空行
|
||||
|
||||
SELECT COALESCE(sum(h.usage_kwh), 0) INTO c0
|
||||
FROM csg_history h
|
||||
WHERE h.kind='day'
|
||||
AND h.period >= date_trunc('month', dd)::date
|
||||
AND h.period < dd;
|
||||
c1 := c0 + v_usage;
|
||||
v_ladder := CASE
|
||||
WHEN c1 <= (CASE WHEN extract(month FROM dd) BETWEEN 5 AND 10 THEN 260 ELSE 200 END) THEN '一档'
|
||||
WHEN c1 <= (CASE WHEN extract(month FROM dd) BETWEEN 5 AND 10 THEN 600 ELSE 400 END) THEN '二档'
|
||||
ELSE '三档' END;
|
||||
|
||||
INSERT INTO csg_history (period, kind, usage_kwh, cost, ladder)
|
||||
VALUES (dd, 'day', v_usage,
|
||||
round(public.csg_ladder_cost_raw(c1, dd) - public.csg_ladder_cost_raw(c0, dd), 2),
|
||||
v_ladder)
|
||||
ON CONFLICT (period, kind) DO NOTHING;
|
||||
END LOOP;
|
||||
|
||||
-- 4b. 受影响月份的 month 行重算(只升不降,与主函数同语义)
|
||||
FOR mm IN
|
||||
SELECT DISTINCT date_trunc('month', h.period)::date
|
||||
FROM csg_history h
|
||||
WHERE h.kind='day' AND h.period >= CURRENT_DATE - lookback
|
||||
ORDER BY 1
|
||||
LOOP
|
||||
SELECT sum(h.usage_kwh) INTO c1
|
||||
FROM csg_history h
|
||||
WHERE h.kind='day' AND h.period >= mm AND h.period < (mm + interval '1 month');
|
||||
SELECT h.usage_kwh INTO c0 FROM csg_history h WHERE h.kind='month' AND h.period = mm;
|
||||
IF c1 IS NOT NULL AND c1 > 0 AND (c0 IS NULL OR c1 > c0) THEN
|
||||
INSERT INTO csg_history (period, kind, usage_kwh, cost)
|
||||
VALUES (mm, 'month', c1, public.csg_ladder_cost(c1, mm))
|
||||
ON CONFLICT (period, kind) DO UPDATE SET
|
||||
usage_kwh = EXCLUDED.usage_kwh,
|
||||
cost = EXCLUDED.cost,
|
||||
updated_at = now();
|
||||
END IF;
|
||||
END LOOP;
|
||||
END $function$;
|
||||
|
||||
-- 重挂自愈任务:每天 15:10 UTC = 23:10 Asia/Shanghai(排在 14:30 UTC 的主任务之后)
|
||||
DO $do$
|
||||
BEGIN
|
||||
IF NOT EXISTS (SELECT 1 FROM timescaledb_information.jobs WHERE proc_name = 'csg_backfill_missing_days') THEN
|
||||
PERFORM add_job(proc => 'public.csg_backfill_missing_days'::regproc,
|
||||
schedule_interval => INTERVAL '1 day',
|
||||
initial_start => date_trunc('day', now() AT TIME ZONE 'UTC')
|
||||
+ INTERVAL '1 day 15 hours 10 minutes',
|
||||
timezone => 'UTC');
|
||||
END IF;
|
||||
END
|
||||
$do$;
|
||||
|
||||
COMMIT;
|
||||
Reference in New Issue
Block a user