性能优化: 任务列表瘦身+详情分页+auto状态独立存储+磁盘统计缓存+db异步同步; 统计数值未读取时显示-
- /api/tasks 响应 1.57MB -> 12KB (latest_run 只带摘要, 一次读 runs.json) - /api/stats 磁盘占用加 30s 缓存, 去除重复全量读 - 详情接口分页返回 results (默认100条/页), 前端表格分页+页码跳转 - auto 任务 pending/visited 队列迁移到 data/auto_state/ 独立文件 (tasks.json 1.6MB -> 14KB) - MySQL 同步改后台线程异步执行 (db.sync_run_async/upsert_task_async), 不再阻塞爬虫和 API - 启动时自动迁移存量 auto 状态 - 统计/回收站/运行中数值在未读到真实数据前显示 '-'
This commit is contained in:
@@ -5,6 +5,7 @@
|
||||
"""
|
||||
import os
|
||||
import threading
|
||||
import time
|
||||
from datetime import datetime
|
||||
|
||||
from flask import Flask, jsonify, request, send_file, send_from_directory
|
||||
@@ -95,7 +96,7 @@ def persist_cb(task_id, run):
|
||||
run["progress"]["percent"] = round(done * 100 / total) if total else 0
|
||||
store.save_run(task_id, run)
|
||||
try:
|
||||
db.sync_run(run, lambda tid, r: store.save_run(tid, r))
|
||||
db.sync_run_async(run) # 后台异步同步, 不再阻塞爬虫线程
|
||||
except Exception as e:
|
||||
print(f"[db] 同步失败: {e}", flush=True)
|
||||
|
||||
@@ -111,7 +112,7 @@ def start_run(task):
|
||||
job = CrawlJob(task, run, persist_cb)
|
||||
JOBS[task["id"]] = job
|
||||
job.start()
|
||||
db.upsert_task(task) # 确保任务在库中
|
||||
db.upsert_task_async(task) # 确保任务在库中 (异步)
|
||||
return run, None
|
||||
|
||||
|
||||
@@ -147,12 +148,37 @@ def api_status():
|
||||
|
||||
# ---------------- API: 任务 ----------------
|
||||
|
||||
def _strip_auto_state(task):
|
||||
"""剥离任务对象中的待爬/已爬队列 (已独立存储, 避免大 JSON 拖慢接口)"""
|
||||
if task.get("mode") == "auto" and task.get("auto"):
|
||||
task["auto"] = {k: v for k, v in task["auto"].items()
|
||||
if k not in ("pending", "visited")}
|
||||
return task
|
||||
|
||||
|
||||
def _run_summary(run):
|
||||
"""运行记录摘要 (不含 results/logs), 供列表/详情页使用"""
|
||||
return {
|
||||
"id": run["id"],
|
||||
"task_id": run.get("task_id", ""),
|
||||
"mode": run.get("mode", ""),
|
||||
"status": run.get("status", ""),
|
||||
"progress": run.get("progress", {}),
|
||||
"started_at": run.get("started_at", ""),
|
||||
"finished_at": run.get("finished_at", ""),
|
||||
"stats": run.get("stats", {}),
|
||||
"out_dir": run.get("out_dir", ""),
|
||||
}
|
||||
|
||||
|
||||
@app.route("/api/tasks", methods=["GET"])
|
||||
def api_tasks():
|
||||
tasks = [t for t in store.load_tasks() if not t.get("deleted_at")]
|
||||
runs_map = store.load_runs_map() # 只全量读一次
|
||||
for t in tasks:
|
||||
runs = store.get_runs(t["id"])
|
||||
t["latest_run"] = runs[-1] if runs else None
|
||||
_strip_auto_state(t)
|
||||
runs = runs_map.get(t["id"], [])
|
||||
t["latest_run"] = _run_summary(runs[-1]) if runs else None
|
||||
with JOBS_LOCK:
|
||||
job = JOBS.get(t["id"])
|
||||
t["running"] = bool(job and job.is_running())
|
||||
@@ -220,7 +246,7 @@ def api_create_task():
|
||||
return jsonify({"error": "请至少填写一个网址"}), 400
|
||||
|
||||
store.upsert_task(task)
|
||||
db.upsert_task(task)
|
||||
db.upsert_task_async(task)
|
||||
return jsonify(task), 201
|
||||
|
||||
|
||||
@@ -229,11 +255,59 @@ def api_task_detail(tid):
|
||||
task = store.get_task(tid)
|
||||
if not task:
|
||||
return jsonify({"error": "任务不存在"}), 404
|
||||
task["runs"] = list(reversed(store.get_runs(tid)))
|
||||
# 分页参数: 指定 run + 页码 (默认最新 run 第 1 页)
|
||||
rid = request.args.get("run", "")
|
||||
try:
|
||||
page = max(int(request.args.get("page", 1)), 1)
|
||||
except ValueError:
|
||||
page = 1
|
||||
try:
|
||||
page_size = min(max(int(request.args.get("page_size", 100)), 10), 500)
|
||||
except ValueError:
|
||||
page_size = 100
|
||||
|
||||
runs = store.get_runs(tid)
|
||||
if rid:
|
||||
cur = next((r for r in runs if r["id"] == rid), None)
|
||||
if cur is None:
|
||||
cur = runs[-1] if runs else None # 指定的 run 不存在时回退最新
|
||||
else:
|
||||
cur = runs[-1] if runs else None
|
||||
with JOBS_LOCK:
|
||||
job = JOBS.get(tid)
|
||||
task["running"] = bool(job and job.is_running())
|
||||
return jsonify(task)
|
||||
|
||||
result = _strip_auto_state(dict(task))
|
||||
result["runs"] = [_run_summary(r) for r in reversed(runs)]
|
||||
result["run"] = None
|
||||
result["run_id"] = ""
|
||||
result["run_total"] = 0
|
||||
result["run_page"] = 1
|
||||
result["run_pages"] = 1
|
||||
if cur:
|
||||
results = cur.get("results", [])
|
||||
total = len(results)
|
||||
pages = max((total + page_size - 1) // page_size, 1)
|
||||
page = min(page, pages)
|
||||
start = (page - 1) * page_size
|
||||
snap = {k: v for k, v in cur.items() if k != "logs"} # logs 走独立轮询接口
|
||||
snap["results"] = results[start:start + page_size]
|
||||
snap["results_total"] = total
|
||||
snap["run_page"] = page
|
||||
snap["run_pages"] = pages
|
||||
snap["run_page_size"] = page_size
|
||||
result["run"] = snap
|
||||
result["run_id"] = cur["id"]
|
||||
result["run_total"] = total
|
||||
result["run_page"] = page
|
||||
result["run_pages"] = pages
|
||||
result["run_page_size"] = page_size
|
||||
# auto 状态数量 (独立文件)
|
||||
if task.get("mode") == "auto":
|
||||
st = store.load_auto_state(tid, task)
|
||||
result["auto_pending_count"] = len(st.get("pending", []))
|
||||
result["auto_visited_count"] = len(st.get("visited", []))
|
||||
return jsonify(result)
|
||||
|
||||
|
||||
@app.route("/api/tasks/<tid>", methods=["PUT"])
|
||||
@@ -267,7 +341,7 @@ def api_update_task(tid):
|
||||
task["schedule"] = sch
|
||||
task["updated_at"] = now_str()
|
||||
store.upsert_task(task)
|
||||
db.upsert_task(task)
|
||||
db.upsert_task_async(task)
|
||||
return jsonify(task)
|
||||
|
||||
|
||||
@@ -283,7 +357,7 @@ def api_delete_task(tid):
|
||||
with JOBS_LOCK:
|
||||
JOBS.pop(tid, None)
|
||||
store.soft_delete_task(tid, now_str())
|
||||
db.upsert_task(store.get_task(tid))
|
||||
db.upsert_task_async(store.get_task(tid))
|
||||
return jsonify({"ok": True, "msg": "已移入回收站"})
|
||||
|
||||
|
||||
@@ -306,8 +380,9 @@ def _purge_out_dir(out_dir):
|
||||
@app.route("/api/trash", methods=["GET"])
|
||||
def api_trash_list():
|
||||
items = store.list_trash()
|
||||
runs_map = store.load_runs_map()
|
||||
for t in items:
|
||||
t["runs_count"] = len(store.get_runs(t["id"]))
|
||||
t["runs_count"] = len(runs_map.get(t["id"], []))
|
||||
t["out_dir"] = resolve_out_dir(t)
|
||||
return jsonify(items)
|
||||
|
||||
@@ -318,7 +393,7 @@ def api_trash_restore(tid):
|
||||
if not task or not task.get("deleted_at"):
|
||||
return jsonify({"error": "任务不在回收站中"}), 404
|
||||
store.restore_task(tid)
|
||||
db.upsert_task(store.get_task(tid))
|
||||
db.upsert_task_async(store.get_task(tid))
|
||||
return jsonify({"ok": True, "msg": "已恢复"})
|
||||
|
||||
|
||||
@@ -386,13 +461,20 @@ def api_resume(tid):
|
||||
|
||||
# ---------------- API: 统计 ----------------
|
||||
|
||||
# 磁盘占用统计缓存 (避免每次轮询都逐个 stat 文件)
|
||||
_DISK_CACHE = {"ts": 0.0, "mb": 0.0}
|
||||
DISK_CACHE_TTL = 30 # 秒
|
||||
|
||||
|
||||
@app.route("/api/stats")
|
||||
def api_stats():
|
||||
tasks = [t for t in store.load_tasks() if not t.get("deleted_at")]
|
||||
trash_count = sum(1 for t in store.load_tasks() if t.get("deleted_at"))
|
||||
tasks = store.load_tasks()
|
||||
active = [t for t in tasks if not t.get("deleted_at")]
|
||||
trash_count = len(tasks) - len(active)
|
||||
runs_map = store.load_runs_map() # 只全量读一次
|
||||
total_runs = ok = fail = imgs = 0
|
||||
for t in tasks:
|
||||
for r in store.get_runs(t["id"]):
|
||||
for t in active:
|
||||
for r in runs_map.get(t["id"], []):
|
||||
total_runs += 1
|
||||
st = r.get("stats") or {}
|
||||
ok += st.get("ok", 0)
|
||||
@@ -400,28 +482,30 @@ def api_stats():
|
||||
imgs += st.get("images", 0)
|
||||
with JOBS_LOCK:
|
||||
running = sum(1 for j in JOBS.values() if j.is_running())
|
||||
# 统计各任务输出目录的磁盘占用
|
||||
size = 0
|
||||
seen = set()
|
||||
for t in tasks:
|
||||
d = os.path.realpath(resolve_out_dir(t))
|
||||
if d in seen or not os.path.isdir(d):
|
||||
continue
|
||||
seen.add(d)
|
||||
for root, _dirs, files in os.walk(d):
|
||||
for f in files:
|
||||
try:
|
||||
size += os.path.getsize(os.path.join(root, f))
|
||||
except OSError:
|
||||
pass
|
||||
now = time.time()
|
||||
if now - _DISK_CACHE["ts"] > DISK_CACHE_TTL:
|
||||
size = 0
|
||||
seen = set()
|
||||
for t in active:
|
||||
d = os.path.realpath(resolve_out_dir(t))
|
||||
if d in seen or not os.path.isdir(d):
|
||||
continue
|
||||
seen.add(d)
|
||||
for root, _dirs, files in os.walk(d):
|
||||
for f in files:
|
||||
try:
|
||||
size += os.path.getsize(os.path.join(root, f))
|
||||
except OSError:
|
||||
pass
|
||||
_DISK_CACHE.update(ts=now, mb=round(size / 1048576, 1))
|
||||
return jsonify({
|
||||
"tasks": len(tasks),
|
||||
"tasks": len(active),
|
||||
"running": running,
|
||||
"runs": total_runs,
|
||||
"ok": ok,
|
||||
"fail": fail,
|
||||
"images": imgs,
|
||||
"disk_mb": round(size / 1048576, 1),
|
||||
"disk_mb": _DISK_CACHE["mb"],
|
||||
"trash": trash_count,
|
||||
})
|
||||
|
||||
@@ -463,9 +547,10 @@ def api_clear_cache(tid):
|
||||
auto = task.setdefault("auto", {})
|
||||
auto["pending"] = []
|
||||
auto["visited"] = []
|
||||
store.clear_auto_state(tid) # 同时清空独立状态文件
|
||||
task["updated_at"] = now_str()
|
||||
store.upsert_task(task)
|
||||
db.upsert_task(task)
|
||||
db.upsert_task_async(task)
|
||||
return jsonify({"ok": True, "msg": "缓存队列已清空"})
|
||||
|
||||
|
||||
@@ -500,6 +585,7 @@ def api_search():
|
||||
ql = q.lower()
|
||||
results = []
|
||||
seen = set()
|
||||
runs_map = store.load_runs_map() # 只全量读一次
|
||||
for task in store.load_tasks():
|
||||
if task.get("deleted_at"):
|
||||
continue # 回收站任务不参与搜索
|
||||
@@ -509,7 +595,7 @@ def api_search():
|
||||
or any(ql in u.lower() for u in task.get("urls", []))
|
||||
or ql in ((auto.get("seed_url") or "").lower())
|
||||
)
|
||||
for run in store.get_runs(task["id"]):
|
||||
for run in runs_map.get(task["id"], []):
|
||||
for r in run.get("results", []):
|
||||
url = r.get("url") or ""
|
||||
title = r.get("title") or ""
|
||||
@@ -593,13 +679,35 @@ def api_file():
|
||||
scheduler = Scheduler(start_run)
|
||||
|
||||
|
||||
def migrate_auto_state():
|
||||
"""启动时把旧任务对象里的 auto pending/visited 迁移到独立文件, 给 tasks.json 瘦身"""
|
||||
migrated = 0
|
||||
for t in store.load_tasks():
|
||||
if t.get("mode") != "auto":
|
||||
continue
|
||||
auto = t.get("auto") or {}
|
||||
if auto.get("pending") or auto.get("visited"):
|
||||
store.save_auto_state(t["id"], {
|
||||
"pending": auto.get("pending", []),
|
||||
"visited": auto.get("visited", []),
|
||||
})
|
||||
auto.pop("pending", None)
|
||||
auto.pop("visited", None)
|
||||
t["updated_at"] = now_str()
|
||||
store.upsert_task(t)
|
||||
migrated += 1
|
||||
if migrated:
|
||||
print(f"[universal-crawler] 已迁移 {migrated} 个 auto 任务的队列状态到独立文件", flush=True)
|
||||
|
||||
|
||||
def main():
|
||||
os.makedirs(os.path.join(HERE, "data"), exist_ok=True)
|
||||
os.makedirs(os.path.join(HERE, "out"), exist_ok=True)
|
||||
db.init_db()
|
||||
# 全量同步存量任务到数据库
|
||||
migrate_auto_state()
|
||||
# 全量同步存量任务到数据库 (异步, 不阻塞启动)
|
||||
for t in store.load_tasks():
|
||||
db.upsert_task(t)
|
||||
db.upsert_task_async(t)
|
||||
scheduler.start()
|
||||
print(f"[universal-crawler] 启动完成, 管理界面: http://0.0.0.0:{PORT}/")
|
||||
app.run(host="0.0.0.0", port=PORT, threaded=True, debug=False)
|
||||
|
||||
Reference in New Issue
Block a user