tradingSystem/app/repo/consensus_stat_repo.py

77 lines
4.0 KiB
Python
Raw Normal View History

# -*- coding: utf-8 -*-
"""pms_consensus_stat 单表访问 (2026-09-14 观察读数包, 台账 013)。
三源合议观察读数的落表与回看**严格单表访问**: 每个函数只碰 pms_consensus_stat 一张表
落表走 execute_many + ON DUPLICATE KEY UPDATE 同一 (北京日期, 代码, 类别) 重复记到时
只把 rounds 加一刷新 reason/detail/last_at, 不新增行 (一天一票一类一行)
回看按 (stat_date, kind) 索引取某段日期的行, 供复核脚本与日报聚合
"""
from __future__ import annotations
from app.db.session import execute, execute_many, fetch_all
# 落表的业务列 (与 ddl_pms_v1.sql 的 pms_consensus_stat 对齐)。
# rounds 不在插入列里 —— 插入时用 DDL 的 DEFAULT 1, 重复时在更新子句里 +1。
_COLS = ("stat_date", "ts_code", "kind", "reason", "detail", "first_at", "last_at")
def upsert_many(rows: list) -> int:
"""批量落观察读数。rows 每项是 {列名: 值}, 至少含 stat_date / ts_code / kind。
(, , ) 重复即累加次数 (rounds+1), 不新增行返回受影响行数"""
rows = [r for r in (rows or [])
if r.get("stat_date") and r.get("ts_code") and r.get("kind")]
if not rows:
return 0
payload = [{c: r.get(c) for c in _COLS} for r in rows]
cols = ", ".join(_COLS)
vals = ", ".join(f":{c}" for c in _COLS)
# ON DUPLICATE 子句只用 VALUES(列), **绝不带绑定参数 (:col)** —— 2026-09-11 真机踩过:
# pymysql 的 executemany 对 INSERT ... ON DUPLICATE 做多行合并, 只展开 VALUES 子句的
# 占位符, UPDATE 子句里的 :col 不展开却仍算参数, 批量时参数错位报 1064。
# first_at 有意不进更新子句 —— 首次记到的时刻要保住; rounds 在库里自增, 不从外面传。
updates = ("reason = VALUES(reason), detail = VALUES(detail), "
"last_at = VALUES(last_at), rounds = rounds + 1")
return execute_many(
f"INSERT INTO pms_consensus_stat ({cols}) VALUES ({vals}) "
f"ON DUPLICATE KEY UPDATE {updates}", payload)
def list_range(date_from: int, date_to: int, kinds=None) -> list:
"""取 [date_from, date_to] (含两端) 的观察读数行, 升序。kinds 非空时只取这些类别。
kinds IN 占位符手动展开 ( tech_repo.history_multi)"""
p = {"a": int(date_from), "b": int(date_to)}
where = "stat_date >= :a AND stat_date <= :b"
ks = [k for k in (kinds or []) if k]
if ks:
keys = []
for i, k in enumerate(ks):
keys.append(f":k{i}")
p[f"k{i}"] = k
where += f" AND kind IN ({', '.join(keys)})"
return fetch_all(
f"SELECT * FROM pms_consensus_stat WHERE {where} "
f"ORDER BY stat_date ASC, kind ASC, ts_code ASC", p)
def delete_day_kind(stat_date: int, kind: str) -> int:
"""删某北京日期某类别的全部行, 返回删除行数。映射重建时先清当日 tech_noread ——
早上两次重建 (06:30/08:40) 各写一遍逐只无读数, 06:30 缺读数08:40 补上的票那一行会留着,
先清再写即以最新一次为准 (2026-09-14 评审第三条)"""
return execute("DELETE FROM pms_consensus_stat WHERE stat_date = :d AND kind = :k",
{"d": int(stat_date), "k": str(kind)})
def count_by_kind(date_from: int, date_to: int) -> dict:
"""[date_from, date_to] 内按 (日期, 类别) 的计数, 返回 {stat_date: {kind: {rows, rounds}}}。
rows 是行数 (几只票), rounds 是累计记到的次数 (同一票记了几轮)"""
rows = fetch_all(
"SELECT stat_date, kind, COUNT(*) AS n_rows, SUM(rounds) AS n_rounds "
"FROM pms_consensus_stat WHERE stat_date >= :a AND stat_date <= :b "
"GROUP BY stat_date, kind", {"a": int(date_from), "b": int(date_to)})
out: dict = {}
for r in rows:
d = int(r["stat_date"])
out.setdefault(d, {})[r["kind"]] = {
"rows": int(r["n_rows"] or 0), "rounds": int(r["n_rounds"] or 0)}
return out