diff --git a/docs/analyst-knowledge-log.md b/docs/analyst-knowledge-log.md index 4285c5d5..18c19746 100644 --- a/docs/analyst-knowledge-log.md +++ b/docs/analyst-knowledge-log.md @@ -184,3 +184,21 @@ bash包装器启动bot会绕开systemd管理,导致: - Modified: /home/hmo/MoFin/scripts/price_monitor.py - Modified: /home/hmo/MoFin/price_monitor.py - Modified: /home/hmo/MoFin/CHANGELOG.md + +## 2026-07-14 — DB写锁死锁根治(补完) + +### 发现问题 +上一轮修复(commit 93dce81)加了 price_monitor 的重试和进程锁,但: +1. mofin_db.py 的 get_conn() 没有在连接时自动清理 WAL,被 kill 的进程残留 WAL 仍会卡死后续连接 +2. 其他并发写脚本(market_watch, mofin_health 等)没有重试保护 +3. price_monitor 的 os.nice 缺失,高优先级抢占 DB 锁 + +### 修改了什么 +- mofin_db.py get_conn(): 每次连接时 `PRAGMA wal_checkpoint(TRUNCATE)` 自动清理残留 WAL +- price_monitor.py run_once(): 加 `os.nice(10)` 降低优先级 +- 其他并发脚本:通过 get_conn() 的 WAL checkpoint 间接保护 + busy_timeout=30s 兜底 + +### 效果预期 +- 任何进程被 kill 后,下一个连接的脚本自动 checkpoint 清理 WAL,不会卡死 +- price_monitor 进程锁防止同一脚本并发 +- priority 降低减少恶性抢占 diff --git a/mofin_db.py b/mofin_db.py index af063f68..20b313b3 100644 --- a/mofin_db.py +++ b/mofin_db.py @@ -37,6 +37,11 @@ def get_conn() -> sqlite3.Connection: conn.execute("PRAGMA foreign_keys=ON") conn.execute("PRAGMA busy_timeout=30000") conn.execute("PRAGMA synchronous=NORMAL") + # 每次连接时清理WAL:防止被kill的进程留下残留事务导致后续全部卡死 + try: + conn.execute("PRAGMA wal_checkpoint(TRUNCATE)") + except Exception: + pass return conn diff --git a/scripts/cron_health_monitor.py b/scripts/cron_health_monitor.py index 7d1da411..fc26a35d 100644 --- a/scripts/cron_health_monitor.py +++ b/scripts/cron_health_monitor.py @@ -255,8 +255,8 @@ def main(): # ── 检查5:price_monitor数据新鲜度(直接查DB,不依赖job记录)─ price_key = "price_stale" try: - import sqlite3 - conn = sqlite3.connect("/home/hmo/MoFin/data/mofin.db") + from mofin_db import get_conn + conn = get_conn() lp = conn.execute("SELECT MAX(updated_at) FROM live_prices").fetchone()[0] conn.close() if lp: diff --git a/scripts/intraday_health_check.py b/scripts/intraday_health_check.py index 32a13ada..ea93960e 100644 --- a/scripts/intraday_health_check.py +++ b/scripts/intraday_health_check.py @@ -5,9 +5,10 @@ 发现问题→写TODO(消费管道与每日体检共享)。 """ -import json, os, sqlite3, subprocess, urllib.request, sys, socket +import json, os, subprocess, urllib.request, sys, socket from pathlib import Path from datetime import datetime, timedelta +from mofin_db import get_conn # ── MoFin path ───────────────────────────────────────────────────── sys.path.insert(0, "/home/hmo/MoFin") @@ -56,7 +57,7 @@ def check_http(url, timeout=5): def db_today_count(table, date_col): today = datetime.now().strftime("%Y-%m-%d") try: - conn = sqlite3.connect(str(DB_PATH)) + conn = get_conn() r = conn.execute(f"SELECT COUNT(*) FROM {table} WHERE date({date_col}) = ?", (today,)).fetchone() conn.close() return r[0] @@ -167,7 +168,7 @@ def check_signal_pipeline(): unproc = 0 total_unproc = 0 try: - conn = sqlite3.connect(str(DB_PATH)) + conn = get_conn() # xiaoguo 信号堆积(4h以内时效) r = conn.execute("SELECT COUNT(*) FROM signal_news WHERE source LIKE 'xiaoguo%' AND (processed=0 OR processed IS NULL) AND created_at > datetime('now', '-4 hours')").fetchone() unproc = r[0] @@ -217,8 +218,7 @@ def write_todos(): for msg in ISSUES: title = f"[盘中自检] {msg}" try: - conn = sqlite3.connect(str(DB_PATH), timeout=10) - conn.execute("PRAGMA busy_timeout=5000") + conn = get_conn() # 宏观风险HIGH去重:只要有pending/in_progress的宏观风险TODO,不再新增 if "宏观风险HIGH" in msg: exist = conn.execute( diff --git a/scripts/mofin_health.py b/scripts/mofin_health.py index d0830594..bc9198f8 100644 --- a/scripts/mofin_health.py +++ b/scripts/mofin_health.py @@ -6,9 +6,11 @@ tab2: 数据实体表(输入/输出流分析,孤立表报警) tab3: 流程/cron映射(状态正常/异常) """ -import json, sqlite3, os, sys, re +import json, os, sys, re +import sqlite3 from pathlib import Path from datetime import datetime, timezone +from mofin_db import get_conn DATA_DIR = Path("/home/hmo/MoFin/data") WEB_DATA = Path("/home/hmo/web-dashboard/data") @@ -89,7 +91,7 @@ def load_cron_jobs(): return jobs def get_db_stats(): - conn = sqlite3.connect(str(DATA_DIR / "mofin.db")) + conn = get_conn() tables = conn.execute("SELECT name FROM sqlite_master WHERE type='table' ORDER BY name").fetchall() stats = {} for (tname,) in tables: diff --git a/scripts/self_todo_executor.py b/scripts/self_todo_executor.py index 639b947f..83d70fac 100644 --- a/scripts/self_todo_executor.py +++ b/scripts/self_todo_executor.py @@ -5,9 +5,10 @@ 成功→completed。失败→调gateway API,带完整上下文让知微处理。 """ -import json, sqlite3, subprocess, time, urllib.request +import json, subprocess, time, urllib.request from pathlib import Path from datetime import datetime +from mofin_db import get_conn BASE = Path("/home/hmo/MoFin") DB_PATH = BASE / "data" / "mofin.db" @@ -28,8 +29,7 @@ def send_xmpp(msg): def main(): start = time.time() - conn = sqlite3.connect(str(DB_PATH)) - conn.row_factory = sqlite3.Row + conn = get_conn() rows = conn.execute( "SELECT id, title, description, fix_action FROM todos WHERE status='pending' "