69 lines
3.5 KiB
Python
69 lines
3.5 KiB
Python
|
|
# -*- 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_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 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
|