修复初始化建表语句

This commit is contained in:
zlt 2026-07-29 09:14:38 +08:00
parent 3a1191660b
commit de6e85cb0b
2 changed files with 98 additions and 25 deletions

View File

@ -79,6 +79,7 @@ class WsRunner:
self._pending_seq = set() # 乱序暂存 (正常恒空) self._pending_seq = set() # 乱序暂存 (正常恒空)
self._unacked = 0 # 距上次 ack 又落了几条 self._unacked = 0 # 距上次 ack 又落了几条
self._baselined = False # 是否已对齐对端序号起点 (见 _set_baseline) self._baselined = False # 是否已对齐对端序号起点 (见 _set_baseline)
self._db_ready = False # 通道三表是否可用 (缺表时空转重试, 不写心跳)
self._params = {} self._params = {}
self._stat = {"rx": 0, "tx": 0, "trades": 0, "rejects": 0, "reconnects": 0, self._stat = {"rx": 0, "tx": 0, "trades": 0, "rejects": 0, "reconnects": 0,
"last_rx_at": None, "last_tx_at": None} "last_rx_at": None, "last_tx_at": None}
@ -111,7 +112,7 @@ class WsRunner:
"idle_timeout_sec": 15, "ack_batch": 20, "idle_timeout_sec": 15, "ack_batch": 20,
"ack_interval_sec": 2.0, "outbox_poll_sec": 0.5, "ack_interval_sec": 2.0, "outbox_poll_sec": 0.5,
"beat_sec": 2, "max_attempts": 3, "connect_timeout": 10} "beat_sec": 2, "max_attempts": 3, "connect_timeout": 10}
logger.warning("参数刷新失败, 沿用上一份: %s", e) logger.warning("参数刷新失败, 沿用上一份: %s", _brief_err(e))
@staticmethod @staticmethod
def _secrets() -> tuple: def _secrets() -> tuple:
@ -127,7 +128,8 @@ class WsRunner:
loop.add_signal_handler(sig, self._request_stop, sig) loop.add_signal_handler(sig, self._request_stop, sig)
await self._refresh_params() await self._refresh_params()
await self._boot() # 注意: 不在这里做 _boot —— 建表检查放进连接循环, 通道没启用时就完全不碰库,
# 表建好之后也能自己恢复, 不必重启容器。
beat = asyncio.create_task(self._beat_loop(), name="beat") beat = asyncio.create_task(self._beat_loop(), name="beat")
try: try:
await self._connection_loop() await self._connection_loop()
@ -142,10 +144,31 @@ class WsRunner:
getattr(sig, "name", sig)) getattr(sig, "name", sig))
self._stop.set() self._stop.set()
async def _boot(self): async def _boot(self) -> bool:
"""启动自检: 恢复水位 + 把上次崩在半路的 SENDING 退回队列。""" """启动自检: 三表在不在 → 恢复水位 → 把上次崩在半路的 SENDING 退回队列。
返回 False 表示还不能开工 (通常是表没建), 由连接循环退回空转并定期重试
"""
try: try:
await _db(qmt_repo.ensure_state) await _db(qmt_repo.ensure_state)
except Exception as e:
self._db_ready = False
msg = _brief_err(e)
if _looks_like_missing_table(msg):
# 缺表是**部署少跑了一步**, 不是异常。给一行照着做就能好的话,
# 而不是六十行 SQLAlchemy 堆栈 —— 后者会把真正该看的信息埋掉。
logger.error(
"ws 通道三表还没建 (%s)。先建表再起本进程:\n"
" docker compose run --rm pms-web python scripts/init_db.py --yes\n"
" docker compose run --rm pms-web python scripts/check_db.py\n"
" 建完**不用重启容器**, 本进程每 60 秒自己重试。在此之前不写心跳, "
"dispatcher 会据此拒发指令 —— 这是对的。", msg)
else:
logger.error("通道表自检失败 (每 60 秒重试): %s", msg)
return False
self._db_ready = True
try:
st = await _db(qmt_repo.get_state) st = await _db(qmt_repo.get_state)
stored = int(st.get("last_seq") or 0) stored = int(st.get("last_seq") or 0)
# last_seq 落库是按批刷的, 可能落后于实际已落库的消息 —— 以 inbox 为准重算。 # last_seq 落库是按批刷的, 可能落后于实际已落库的消息 —— 以 inbox 为准重算。
@ -160,12 +183,14 @@ class WsRunner:
# 只回 ack{duplicate:true} 带当前状态。幂等键就是为这一刻准备的。 # 只回 ack{duplicate:true} 带当前状态。幂等键就是为这一刻准备的。
logger.warning("%s 张委托上次卡在 SENDING, 已退回队列重发 (幂等键兜底)", n) logger.warning("%s 张委托上次卡在 SENDING, 已退回队列重发 (幂等键兜底)", n)
except Exception as e: except Exception as e:
logger.exception("启动自检失败: %s", e) logger.error("启动自检失败, 水位按 0 起算 (下轮重连会重试): %s", _brief_err(e))
return False
return True
def _operable(self) -> bool: def _operable(self) -> bool:
"""通道是否**有可能**工作 (已启用 + 密钥齐)。不满足就不写心跳 —— 见 _beat_loop。""" """通道是否**有可能**工作 (已启用 + 密钥齐 + 表已建)。不满足就不写心跳 —— 见 _beat_loop。"""
seed, peer = self._secrets() seed, peer = self._secrets()
return bool(self._p("enabled", False) and seed and peer) return bool(self._p("enabled", False) and seed and peer and self._db_ready)
async def _beat_loop(self): async def _beat_loop(self):
"""存活心跳。**连不上 QMT 也照跳** —— 它证明的是「进程在」而不是「连接通」, """存活心跳。**连不上 QMT 也照跳** —— 它证明的是「进程在」而不是「连接通」,
@ -180,7 +205,7 @@ class WsRunner:
if self._operable(): if self._operable():
await _db(qmt_repo.beat, self._stat) await _db(qmt_repo.beat, self._stat)
except Exception as e: except Exception as e:
logger.warning("心跳写入失败 (库不可用?): %s", e) logger.warning("心跳写入失败 (库不可用?): %s", _brief_err(e))
with contextlib.suppress(asyncio.TimeoutError): with contextlib.suppress(asyncio.TimeoutError):
await asyncio.wait_for(self._stop.wait(), timeout=self._p("beat_sec", 2)) await asyncio.wait_for(self._stop.wait(), timeout=self._p("beat_sec", 2))
@ -194,9 +219,15 @@ class WsRunner:
if not seed or not peer: if not seed or not peer:
await self._idle_note( await self._idle_note(
"缺少 Ed25519 密钥: PMS_QMT_SIGN_SEED_HEX / PMS_QMT_PEER_PUBKEY_B64 " "缺少 Ed25519 密钥: PMS_QMT_SIGN_SEED_HEX / PMS_QMT_PEER_PUBKEY_B64 "
"须在 .env 注入 (协议 §10.1.1)。未配齐前不连接, 也不写心跳 —— " "须在 .env 注入 (协议 §10.1.1)。生成与自检: "
"python scripts/gen_keys.py [--check]。未配齐前不连接, 也不写心跳 —— "
"dispatcher 会因此一律拒发, 这是对的", error=True) "dispatcher 会因此一律拒发, 这是对的", error=True)
continue continue
if not self._db_ready and not await self._boot():
# 建表提示已在 _boot 里打过一次, 这里只安静等 —— 不重复刷屏
with contextlib.suppress(asyncio.TimeoutError):
await asyncio.wait_for(self._stop.wait(), timeout=60)
continue
url = self._p("url", settings.PMS_QMT_WS_URL) url = self._p("url", settings.PMS_QMT_WS_URL)
try: try:
@ -426,7 +457,7 @@ class WsRunner:
pl.get("scope"), pl.get("error")) pl.get("scope"), pl.get("error"))
except Exception as e: except Exception as e:
# 通道状态更新失败不该拖垮连接 —— 账本不靠它, 靠 inbox 里的 trade。 # 通道状态更新失败不该拖垮连接 —— 账本不靠它, 靠 inbox 里的 trade。
logger.warning("通道状态更新失败 %s %s: %s", type_, iid, e) logger.warning("通道状态更新失败 %s %s: %s", type_, iid, _brief_err(e))
async def _on_order_update(self, iid: str, pl: dict): async def _on_order_update(self, iid: str, pl: dict):
status = str(pl.get("status") or "").upper() status = str(pl.get("status") or "").upper()
@ -478,12 +509,12 @@ class WsRunner:
try: try:
await _db(qmt_repo.save_watermark, self._last_seq, self._acked_seq) await _db(qmt_repo.save_watermark, self._last_seq, self._acked_seq)
except Exception as e: except Exception as e:
logger.warning("水位落库失败, 本轮不 ack (宁可让对端多留一会儿): %s", e) logger.warning("水位落库失败, 本轮不 ack (宁可让对端多留一会儿): %s", _brief_err(e))
return return
try: try:
await self._send(wsc.T_ACK_SEQ, wsc.ack_seq_payload(self._last_seq), seed) await self._send(wsc.T_ACK_SEQ, wsc.ack_seq_payload(self._last_seq), seed)
except Exception as e: except Exception as e:
logger.warning("ack_seq 发送失败 (下轮重试): %s", e) logger.warning("ack_seq 发送失败 (下轮重试): %s", _brief_err(e))
return return
self._acked_seq = self._last_seq self._acked_seq = self._last_seq
self._unacked = 0 self._unacked = 0
@ -515,7 +546,7 @@ class WsRunner:
except Exception as e: except Exception as e:
if _is_conn_error(e): if _is_conn_error(e):
raise raise
logger.exception("出口轮询异常 (不产生新指令, 下轮继续): %s", e) logger.exception("出口轮询异常 (不产生新指令, 下轮继续): %s", _brief_err(e))
with contextlib.suppress(asyncio.TimeoutError): with contextlib.suppress(asyncio.TimeoutError):
await asyncio.wait_for(self._stop.wait(), await asyncio.wait_for(self._stop.wait(),
timeout=self._p("outbox_poll_sec", 0.5)) timeout=self._p("outbox_poll_sec", 0.5))
@ -608,6 +639,24 @@ def _ymd() -> int:
return int(datetime.now().strftime("%Y%m%d")) return int(datetime.now().strftime("%Y%m%d"))
def _brief_err(e: BaseException) -> str:
"""异常压成一行。
SQLAlchemy 的报错里会把整条 SQL 和全部参数带上, 动辄十几行; 再叠一层
logger.exception 的堆栈, 一个表还没建能刷出六十行, 真正有用的那半句反而被埋掉
日志是给人看的, 定位靠的是第一句话
"""
first = str(e).splitlines()[0].strip()
return f"{type(e).__name__}: {first[:220]}"
def _looks_like_missing_table(msg: str) -> bool:
"""这条报错是不是「表不存在」。MySQL 是 1146, ShardingSphere 代理回的是 10002。"""
low = msg.lower()
return ("does not exist" in low or "doesn't exist" in low
or "1146" in low or "10002" in low)
def _is_conn_error(e: BaseException) -> bool: def _is_conn_error(e: BaseException) -> bool:
"""这个异常是不是"连接没了" """这个异常是不是"连接没了"

View File

@ -60,7 +60,18 @@ def split_sql(text: str) -> list:
def parse_statements(sql_text: str) -> list: def parse_statements(sql_text: str) -> list:
"""去掉整行 `--` 注释后切语句。返回 [(表名, 语句)]; 非 CREATE TABLE 的残句会被标 '?'""" """去掉整行 `--` 注释后切语句。返回 [(种类, 表名, 语句)]。
种类: `table` = CREATE TABLE / `seed` = 幂等的初始化 INSERT / `?` = 认不出的残句
seed 是因为 `pms_ws_state` **单行表**, 那一行 (id=1) 本身就属于表结构的一部分,
跟建表放一起最省事但只放行 `INSERT ... ON DUPLICATE KEY UPDATE` / `INSERT IGNORE`
这种可重复执行的写法 建表脚本必须幂等, 一条会重复插数的 INSERT 混进来, 比认不出
它更危险
`?` 一律中止且不执行任何语句: 它通常意味着 split_sql 把某条 CREATE TABLE 劈断了
(列注释里的分号是老坑), 这时候硬执行残句只会建出半张表
"""
body = "\n".join(ln for ln in sql_text.splitlines() if not ln.strip().startswith("--")) body = "\n".join(ln for ln in sql_text.splitlines() if not ln.strip().startswith("--"))
out = [] out = []
for stmt in split_sql(body): for stmt in split_sql(body):
@ -68,7 +79,14 @@ def parse_statements(sql_text: str) -> list:
if not s: if not s:
continue continue
m = re.search(r"CREATE\s+TABLE\s+(?:IF\s+NOT\s+EXISTS\s+)?`?([A-Za-z0-9_]+)`?", s, re.I) m = re.search(r"CREATE\s+TABLE\s+(?:IF\s+NOT\s+EXISTS\s+)?`?([A-Za-z0-9_]+)`?", s, re.I)
out.append((m.group(1) if m else "?", s)) if m:
out.append(("table", m.group(1), s))
continue
m = re.match(r"INSERT\s+(?:IGNORE\s+)?INTO\s+`?([A-Za-z0-9_]+)`?", s, re.I)
if m and re.search(r"ON\s+DUPLICATE\s+KEY\s+UPDATE|INSERT\s+IGNORE", s, re.I):
out.append(("seed", m.group(1), s))
continue
out.append(("?", "?", s))
return out return out
@ -85,18 +103,22 @@ def main():
stmts = parse_statements(f.read()) stmts = parse_statements(f.read())
only = {x.strip() for x in args.only.split(",") if x.strip()} only = {x.strip() for x in args.only.split(",") if x.strip()}
if only: if only:
stmts = [(t, s) for t, s in stmts if t in only] stmts = [(k, t, s) for k, t, s in stmts if t in only]
if not stmts: if not stmts:
print("FAIL 没有可执行的语句 (检查 --only 是否写对)") print("FAIL 没有可执行的语句 (检查 --only 是否写对)")
sys.exit(1) sys.exit(1)
broken = [s for t, s in stmts if t == "?"] broken = [s for k, _, s in stmts if k == "?"]
if broken: if broken:
print(f"FAIL 解析出 {len(broken)}非 CREATE TABLE 残句, 说明 DDL 切分有误, " print(f"FAIL 解析出 {len(broken)}既不是 CREATE TABLE 也不是幂等 INSERT 的残句, "
f"已中止 (不执行任何语句)。首条残句:\n{broken[0][:200]}") f"说明 DDL 切分有误, 已中止 (不执行任何语句)。首条残句:\n{broken[0][:200]}")
sys.exit(1) sys.exit(1)
tables = [t for k, t, _ in stmts if k == "table"]
seeds = [t for k, t, _ in stmts if k == "seed"]
print(f"DDL 文件: {DDL_FILE}") print(f"DDL 文件: {DDL_FILE}")
print(f"待建表 {len(stmts)} 张: {', '.join(t for t, _ in stmts)}") print(f"待建表 {len(tables)} 张: {', '.join(tables)}")
if seeds:
print(f"待写初始行 {len(seeds)} 条 (幂等): {', '.join(seeds)}")
if not args.yes: if not args.yes:
print("\n[演练模式] 未执行任何语句。确认无误后加 --yes 重跑:") print("\n[演练模式] 未执行任何语句。确认无误后加 --yes 重跑:")
print(" docker compose run --rm pms-web python scripts/init_db.py --yes") print(" docker compose run --rm pms-web python scripts/init_db.py --yes")
@ -109,19 +131,20 @@ def main():
eng = get_engine("proxy") eng = get_engine("proxy")
okc, failed = 0, [] okc, failed = 0, []
print() print()
for tbl, stmt in stmts: for kind, tbl, stmt in stmts:
label = tbl if kind == "table" else f"{tbl} (初始行)"
try: try:
with eng.begin() as c: with eng.begin() as c:
c.execute(text(stmt)) c.execute(text(stmt))
print(f" OK {tbl}") print(f" OK {label}")
okc += 1 okc += 1
except Exception as e: except Exception as e:
print(f" FAIL {tbl}: {type(e).__name__}: {e}") print(f" FAIL {label}: {type(e).__name__}: {e}")
failed.append((tbl, stmt, str(e))) failed.append((tbl, stmt, str(e)))
print("\n[验证] 逐表 COUNT(*)") print("\n[验证] 逐表 COUNT(*)")
missing = [] missing = []
for tbl, _ in stmts: for tbl in dict.fromkeys(t for _, t, _ in stmts): # 去重且保序
try: try:
with eng.connect() as c: with eng.connect() as c:
n = c.execute(text(f"SELECT COUNT(*) FROM {tbl}")).scalar() n = c.execute(text(f"SELECT COUNT(*) FROM {tbl}")).scalar()
@ -138,7 +161,8 @@ def main():
if missing: if missing:
print(f"\nFAILED: {len(missing)} 张表仍不可用: {', '.join(missing)}") print(f"\nFAILED: {len(missing)} 张表仍不可用: {', '.join(missing)}")
sys.exit(1) sys.exit(1)
print(f"ALL OK: {okc} 张表就绪, 接着跑 python scripts/check_db.py 做完整自检") print(f"ALL OK: {len(tables)} 张表就绪 ({okc} 条语句执行成功), "
f"接着跑 python scripts/check_db.py 做完整自检")
if __name__ == "__main__": if __name__ == "__main__":