diff --git a/README.md b/README.md index 9365f55..b9d7e6f 100644 --- a/README.md +++ b/README.md @@ -63,6 +63,13 @@ - **图片集**(`<文件名>_img/meta.json`):所属页面、来源链接、每张图片的原始 URL / 大小 / 下载时间 - 详情页结果表中点「📋 元数据」即可在线查看 +### 6. 数据库记录(MySQL) +任务、运行记录、爬取结果实时写入 MySQL(`121.40.164.32:16006`,库 `crawler`),**网页完整内容不入库**(存磁盘文件),库中只存标题、网址、状态、文件路径等元数据: +- `crawl_tasks` — 任务信息(含回收站标记 deleted_at) +- `crawl_runs` — 每次运行记录(状态/进度/成功失败数/图片数/时间) +- `crawl_results` — 每页一条(**status: OK/FAIL** 成功失败标记、标题、网址、来源链接、深度、错误信息、文件路径、图片数) +- 数据库不可用时自动降级,不影响爬取主流程;删除任务进回收站同步标记,彻底删除同步清库 + ## 输出文件 每个任务输出到独立目录(默认 `out/<任务ID>/`): diff --git a/app.py b/app.py index 17aabfe..d0a4721 100644 --- a/app.py +++ b/app.py @@ -10,6 +10,7 @@ from datetime import datetime from flask import Flask, jsonify, request, send_file, send_from_directory import store +import db from engine import CrawlJob, probe_links from scheduler import Scheduler, cron_next, interval_delta @@ -93,6 +94,10 @@ def persist_cb(task_id, run): done = run["progress"].get("done") or 0 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)) + except Exception as e: + print(f"[db] 同步失败: {e}", flush=True) def start_run(task): @@ -106,6 +111,7 @@ def start_run(task): job = CrawlJob(task, run, persist_cb) JOBS[task["id"]] = job job.start() + db.upsert_task(task) # 确保任务在库中 return run, None @@ -214,6 +220,7 @@ def api_create_task(): return jsonify({"error": "请至少填写一个网址"}), 400 store.upsert_task(task) + db.upsert_task(task) return jsonify(task), 201 @@ -260,6 +267,7 @@ def api_update_task(tid): task["schedule"] = sch task["updated_at"] = now_str() store.upsert_task(task) + db.upsert_task(task) return jsonify(task) @@ -275,6 +283,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)) return jsonify({"ok": True, "msg": "已移入回收站"}) @@ -309,6 +318,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)) return jsonify({"ok": True, "msg": "已恢复"}) @@ -320,6 +330,7 @@ def api_trash_purge(tid): out_dir = resolve_out_dir(task) removed = _purge_out_dir(out_dir) store.purge_task(tid) + db.purge_task_db(tid) return jsonify({"ok": True, "purged": True, "files_removed": removed, "out_dir": out_dir}) @@ -328,6 +339,7 @@ def api_trash_clear(): items = store.list_trash() dirs = [resolve_out_dir(t) for t in items] store.purge_trash() + db.purge_trash_db() removed = sum(1 for d in dirs if _purge_out_dir(d)) return jsonify({"ok": True, "purged": len(items), "dirs_removed": removed}) @@ -563,6 +575,10 @@ scheduler = Scheduler(start_run) 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() + # 全量同步存量任务到数据库 + for t in store.load_tasks(): + db.upsert_task(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) diff --git a/db.py b/db.py new file mode 100644 index 0000000..f07c719 --- /dev/null +++ b/db.py @@ -0,0 +1,264 @@ +# -*- coding: utf-8 -*- +""" +MySQL 记录层: 任务 / 运行 / 爬取结果 写入数据库 +- 网页完整内容不入库 (只存磁盘文件), 库中只存标题、网址、状态、文件路径等元数据 +- 爬取成功/失败均有 status 标记 (OK / FAIL), 运行记录有 run 状态标记 +- 所有操作容错: 数据库不可用时不影响爬取主流程 (仅打印日志) +""" +import json +import os +import time + +import pymysql + +DB_CONFIG = dict( + host="121.40.164.32", + port=16006, + user="root", + password="hz_123", + charset="utf8mb4", +) + +DB_NAME = "crawler" + +DDL = [ + """ + CREATE TABLE IF NOT EXISTS crawl_tasks ( + id VARCHAR(32) PRIMARY KEY, + name VARCHAR(200) NOT NULL, + mode VARCHAR(20) NOT NULL, + config TEXT, + urls TEXT, + auto_config TEXT, + schedule_config TEXT, + created_at DATETIME, + updated_at DATETIME, + deleted_at DATETIME NULL + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 + """, + """ + CREATE TABLE IF NOT EXISTS crawl_runs ( + id VARCHAR(32) PRIMARY KEY, + task_id VARCHAR(32) NOT NULL, + mode VARCHAR(20) NOT NULL, + status VARCHAR(20) NOT NULL DEFAULT 'running', + total INT DEFAULT 0, + done INT DEFAULT 0, + ok_count INT DEFAULT 0, + fail_count INT DEFAULT 0, + image_count INT DEFAULT 0, + started_at DATETIME NULL, + finished_at DATETIME NULL, + out_dir VARCHAR(500), + KEY idx_task (task_id) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 + """, + """ + CREATE TABLE IF NOT EXISTS crawl_results ( + id BIGINT AUTO_INCREMENT PRIMARY KEY, + run_id VARCHAR(32) NOT NULL, + task_id VARCHAR(32) NOT NULL, + mode VARCHAR(20), + url VARCHAR(2000) NOT NULL, + title VARCHAR(500), + status VARCHAR(10) NOT NULL, + error TEXT, + crawl_time DATETIME, + source_url VARCHAR(2000), + depth INT, + attempts INT DEFAULT 1, + html_file VARCHAR(500), + txt_file VARCHAR(500), + meta_file VARCHAR(500), + image_count INT DEFAULT 0, + image_files TEXT, + UNIQUE KEY uk_run_url (run_id, url(500)) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 + """, +] + + +def _conn(): + cfg = dict(DB_CONFIG) + cfg["database"] = DB_NAME + return pymysql.connect(**cfg, autocommit=True, connect_timeout=5) + + +def _safe(fn, *args, **kwargs): + """执行数据库操作, 失败仅打印日志不抛出""" + try: + fn(*args, **kwargs) + except Exception as e: + print(f"[db] 操作失败: {e}", flush=True) + + +def init_db(): + """建库建表 (幂等)""" + try: + conn = pymysql.connect(**DB_CONFIG, autocommit=True, connect_timeout=5) + with conn.cursor() as cur: + cur.execute(f"CREATE DATABASE IF NOT EXISTS `{DB_NAME}` DEFAULT CHARACTER SET utf8mb4") + conn.close() + conn = _conn() + with conn.cursor() as cur: + for ddl in DDL: + cur.execute(ddl) + conn.close() + print("[db] 数据库初始化完成 (crawler: crawl_tasks / crawl_runs / crawl_results)") + except Exception as e: + print(f"[db] 数据库初始化失败: {e}", flush=True) + + +# ---------------- 序列化工具 ---------------- + +def _dt(v): + """datetime 字符串 -> MySQL DATETIME (无效返回 None)""" + if not v: + return None + s = str(v).replace("T", " ") + if len(s) >= 19: + return s[:19] + return None + + +def _j(v): + return json.dumps(v, ensure_ascii=False) if v else None + + +# ---------------- 任务 ---------------- + +def upsert_task(task): + """任务写入/更新 (含回收站状态 deleted_at)""" + def _do(): + conn = _conn() + with conn.cursor() as cur: + cur.execute( + """INSERT INTO crawl_tasks + (id, name, mode, config, urls, auto_config, schedule_config, + created_at, updated_at, deleted_at) + VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s) + ON DUPLICATE KEY UPDATE + name=%s, mode=%s, config=%s, urls=%s, auto_config=%s, + schedule_config=%s, updated_at=%s, deleted_at=%s""", + ( + task["id"], task.get("name", ""), task.get("mode", ""), + _j(task.get("config")), _j(task.get("urls")), + _j(task.get("auto")), _j(task.get("schedule")), + _dt(task.get("created_at")), _dt(task.get("updated_at")), + _dt(task.get("deleted_at")), + task.get("name", ""), task.get("mode", ""), + _j(task.get("config")), _j(task.get("urls")), + _j(task.get("auto")), _j(task.get("schedule")), + _dt(task.get("updated_at")), _dt(task.get("deleted_at")), + ), + ) + conn.close() + _safe(_do) + + +def purge_task_db(task_id): + """彻底删除任务记录""" + def _do(): + conn = _conn() + with conn.cursor() as cur: + cur.execute("DELETE FROM crawl_results WHERE task_id=%s", (task_id,)) + cur.execute("DELETE FROM crawl_runs WHERE task_id=%s", (task_id,)) + cur.execute("DELETE FROM crawl_tasks WHERE id=%s", (task_id,)) + conn.close() + _safe(_do) + + +def purge_trash_db(): + """清空回收站 (删除所有已标记删除的任务记录)""" + def _do(): + conn = _conn() + with conn.cursor() as cur: + cur.execute("SELECT id FROM crawl_tasks WHERE deleted_at IS NOT NULL") + ids = [r[0] for r in cur.fetchall()] + for tid in ids: + cur.execute("DELETE FROM crawl_results WHERE task_id=%s", (tid,)) + cur.execute("DELETE FROM crawl_runs WHERE task_id=%s", (tid,)) + cur.execute("DELETE FROM crawl_tasks WHERE id=%s", (tid,)) + conn.close() + _safe(_do) + + +# ---------------- 运行记录 ---------------- + +def upsert_run(run): + """运行记录写入/更新""" + def _do(): + conn = _conn() + st = run.get("stats") or {} + with conn.cursor() as cur: + cur.execute( + """INSERT INTO crawl_runs + (id, task_id, mode, status, total, done, + ok_count, fail_count, image_count, + started_at, finished_at, out_dir) + VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s) + ON DUPLICATE KEY UPDATE + status=%s, total=%s, done=%s, ok_count=%s, fail_count=%s, + image_count=%s, finished_at=%s, out_dir=%s""", + ( + run["id"], run.get("task_id", ""), run.get("mode", ""), + run.get("status", ""), run.get("progress", {}).get("total", 0), + run.get("progress", {}).get("done", 0), + st.get("ok", 0), st.get("fail", 0), st.get("images", 0), + _dt(run.get("started_at")), _dt(run.get("finished_at")), + run.get("out_dir", ""), + run.get("status", ""), run.get("progress", {}).get("total", 0), + run.get("progress", {}).get("done", 0), + st.get("ok", 0), st.get("fail", 0), st.get("images", 0), + _dt(run.get("finished_at")), run.get("out_dir", ""), + ), + ) + conn.close() + _safe(_do) + + +# ---------------- 爬取结果 (每页一条, 成功失败均记录) ---------------- + +def insert_results(run, results): + """批量插入爬取结果 (增量)""" + if not results: + return + + def _do(): + conn = _conn() + rows = [] + for r in results: + rows.append(( + run["id"], run.get("task_id", ""), run.get("mode", ""), + (r.get("url") or "")[:2000], (r.get("title") or "")[:500], + r.get("status", "FAIL"), r.get("error"), + _dt(r.get("crawl_time")), (r.get("source_url") or "")[:2000], + r.get("depth"), r.get("attempts", 1), + r.get("html_file", ""), r.get("txt_file", ""), r.get("meta_file", ""), + len(r.get("images", []) or []), + _j([im.get("file") for im in (r.get("images") or [])]), + )) + with conn.cursor() as cur: + cur.executemany( + """INSERT IGNORE INTO crawl_results + (run_id, task_id, mode, url, title, status, error, + crawl_time, source_url, depth, attempts, + html_file, txt_file, meta_file, image_count, image_files) + VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)""", + rows, + ) + conn.close() + _safe(_do) + + +def sync_run(run, persist): + """持久化回调: 同步运行记录 + 增量同步爬取结果 + persist: callable(task_id, run) 用于回写已同步进度标记 + """ + upsert_run(run) + results = run.get("results", []) + synced = run.get("_db_count", 0) + if len(results) > synced: + insert_results(run, results[synced:]) + run["_db_count"] = len(results) + persist(run.get("task_id"), run) diff --git a/requirements.txt b/requirements.txt index 79264a8..cfbf0a5 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,3 +1,4 @@ flask>=3.0 playwright>=1.40 playwright-stealth>=1.0 +pymysql>=1.1 diff --git a/scheduler.py b/scheduler.py index 438b9b2..6f034e5 100644 --- a/scheduler.py +++ b/scheduler.py @@ -116,5 +116,10 @@ class Scheduler(threading.Thread): sch["last_run"] = now.strftime("%Y-%m-%d %H:%M:%S") sch["runs_count"] = sch.get("runs_count", 0) + 1 store.upsert_task(task) + try: + import db + db.upsert_task(task) + except Exception: + pass except Exception: pass