Compare commits

...
9 Commits
Author SHA1 Message Date
hz4th_coder 66ba1a7cfc v1.0.14: 自动爬取支持深度/页数无限制(默认0/1000), 待爬链接持久化缓存队列(停止后续爬), 起始网址不去重, 清空缓存API 2026-08-11 21:07:31 +08:00
hz4th_coder 3ad3d00ca8 v1.0.13: 自动爬取改为持续递归(无深度限制, 爬到无新链接为止), max_pages作安全上限(默认200) 2026-08-11 19:45:51 +08:00
hz4th_coder aa0c2f56d0 v1.0.12: crawl_results 合并 html_file/txt_file/meta_file 为单一 base_file(同前缀不同后缀), 含存量数据自动迁移 2026-08-11 19:35:55 +08:00
hz4th_coder bf6be47c01 v1.0.11: 新增历史记录回填脚本 backfill.py(幂等, 支持指定任务), 补齐数据库功能上线前的爬取记录 2026-08-11 19:29:57 +08:00
hz4th_coder abde505d32 v1.0.10: 修正数据库连接信息(uni_crawler/wleD6x2T, 库uni_crawler) 2026-08-11 19:27:18 +08:00
hz4th_coder 9cd5b29644 v1.0.9: 爬取记录入库MySQL(crawler库: 任务/运行/结果三表, 内容不入库只存元数据, OK/FAIL状态标记) 2026-08-11 19:23:35 +08:00
hz4th_coder 5a137514e8 v1.0.8: 自动爬取URL规范化去重(锚点/跟踪参数/尾部斜杠/默认端口) + 修复反爬启发式误伤小页面 2026-08-11 17:57:59 +08:00
hz4th_coder 13a760fa6f v1.0.7: 任务列表按模式分组展示(批量/定时/自动各占一行, 空组隐藏) 2026-08-11 15:40:44 +08:00
hz4th_coder add2cea545 v1.0.6: 导入网址文件增加追加/覆盖处理方式选择(默认追加) 2026-08-11 15:33:19 +08:00
10 changed files with 545 additions and 37 deletions
+15 -3
View File
@@ -33,7 +33,10 @@
- 运行中的任务参数支持**热更新**(修改后从下一页起生效)
### 2. 批量爬取模式
一次粘贴多个网址(每行一个,`#` 注释),或点击「📂 导入网址文件」上传 .txt 文件批量导入(自动识别 UTF-8/GBK 编码、自动去重、自动清理行内注释)
一次粘贴多个网址(每行一个,`#` 注释),或点击「📂 导入网址文件」上传 .txt 文件批量导入:
- **处理方式可选**:追加(保留已有,默认)或覆盖(清空已有)
- 自动识别 UTF-8/GBK 编码、自动去重、自动清理行内注释
可配置:
- 项目名称、输出目录(默认 `out/<任务ID>`,可填绝对路径)
- 爬取间隔(随机秒数区间,防封 IP)、单页超时
- 失败重试次数 / 重试间隔
@@ -46,12 +49,14 @@
- 可启用/停用调度,自动计算下次执行时间;到点自动开跑,跑完自动计算下一次
### 4. 自动爬取模式
给定一个起始网址,系统自动从页面里发现链接、按规则筛选后 BFS 爬取:
给定一个起始网址,系统**持续递归**爬取:爬取页面 → 自动发现符合规则的链接 → 继续爬取 → 继续发现……直到**无新链接可爬**时自动结束。
- **🧪 试爬取**:先用起始网址试跑一次,展示规则筛选后的链接清单(将爬取哪些、排除哪些及原因),确认规则符合预期后再正式开爬
- **最大爬取深度**:0=无限制(默认),N=只爬 N 层
- **最大页数**:0=无限制,默认 1000 作安全上限(防止动态无限链接的站点失控)
- **🔄 缓存续爬**:提取但未爬取的链接自动存入缓存队列(连同已爬集合一并持久化),任务停止/中断后再次运行,从缓存队列**继续爬取**,已爬过的不会重复;起始网址每次运行都重新爬取(不去重);规则变更后可点「🧹 清空缓存」重新开始
- **包含规则**:只爬包含指定子串(或正则)的链接
- **排除规则**:跳过匹配的链接(如 login、/tag/
- **仅同域名**:限制在起始网站内
- **最大页数 / 最大深度**:控制爬取规模
- 其余参数(间隔、重试、图片、通知)同批量模式
### 5. 资源操作信息(元数据)
@@ -60,6 +65,13 @@
- **图片集**`<文件名>_img/meta.json`):所属页面、来源链接、每张图片的原始 URL / 大小 / 下载时间
- 详情页结果表中点「📋 元数据」即可在线查看
### 6. 数据库记录(MySQL
任务、运行记录、爬取结果实时写入 MySQL(`121.40.164.32:16006`,账号 `uni_crawler`,库 `uni_crawler`),**网页完整内容不入库**(存磁盘文件),库中只存标题、网址、状态、文件路径等元数据:
- `crawl_tasks` — 任务信息(含回收站标记 deleted_at
- `crawl_runs` — 每次运行记录(状态/进度/成功失败数/图片数/时间)
- `crawl_results` — 每页一条(**status: OK/FAIL** 成功失败标记、标题、网址、来源链接、深度、错误信息、文件路径、图片数)
- 数据库不可用时自动降级,不影响爬取主流程;删除任务进回收站同步标记,彻底删除同步清库
## 输出文件
每个任务输出到独立目录(默认 `out/<任务ID>/`):
+37
View File
@@ -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})
@@ -436,6 +448,27 @@ def api_probe():
return jsonify(result)
@app.route("/api/tasks/<tid>/clear_cache", methods=["POST"])
def api_clear_cache(tid):
"""清空自动任务的待爬缓存队列与已爬集合 (规则变更后重新开始用)"""
task = store.get_task(tid)
if not task:
return jsonify({"error": "任务不存在"}), 404
if task.get("mode") != "auto":
return jsonify({"error": "仅自动爬取任务支持清空缓存"}), 400
with JOBS_LOCK:
job = JOBS.get(tid)
if job and job.is_running():
return jsonify({"error": "任务正在运行,无法清空缓存"}), 409
auto = task.setdefault("auto", {})
auto["pending"] = []
auto["visited"] = []
task["updated_at"] = now_str()
store.upsert_task(task)
db.upsert_task(task)
return jsonify({"ok": True, "msg": "缓存队列已清空"})
@app.route("/api/tasks/<tid>/probe", methods=["POST"])
def api_task_probe(tid):
"""对已保存的自动任务执行试爬取 (使用保存的规则)"""
@@ -563,6 +596,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)
+26
View File
@@ -0,0 +1,26 @@
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
历史爬取记录回填数据库 (幂等, 可重复执行)
用法:
python backfill.py # 回填所有任务的历史记录
python backfill.py <task_id> ... # 只回填指定任务
说明:
- 将本地 data/ 中的任务/运行/爬取结果(含成功与失败)全量写入 MySQL
- 网页内容不入库, 只写元数据; 重复执行不会产生重复记录
"""
import sys
import db
import store
if __name__ == "__main__":
ids = [a for a in sys.argv[1:] if a.strip()] or None
db.init_db()
stats = db.sync_all_history(ids)
if ids:
print(f"[backfill] 已回填 {len(ids)} 个任务: {stats}")
else:
print(f"[backfill] 已回填全部任务: {stats}")
+314
View File
@@ -0,0 +1,314 @@
# -*- 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="uni_crawler",
password="wleD6x2T",
charset="utf8mb4",
)
DB_NAME = "uni_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,
base_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 _migrate(conn):
"""存量表结构迁移: 合并 html_file/txt_file/meta_file 为 base_file"""
with conn.cursor() as cur:
cur.execute("SHOW COLUMNS FROM crawl_results LIKE 'html_file'")
if cur.fetchone():
cur.execute("ALTER TABLE crawl_results ADD COLUMN base_file VARCHAR(500) NULL AFTER meta_file")
cur.execute("UPDATE crawl_results SET base_file = REPLACE(html_file, '.html', '') WHERE base_file IS NULL")
cur.execute("ALTER TABLE crawl_results DROP COLUMN html_file, DROP COLUMN txt_file, DROP COLUMN meta_file")
print("[db] 表结构迁移完成: html_file/txt_file/meta_file -> base_file")
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()
_migrate(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 _base_of(entry):
"""从结果条目提取基础文件名 (html/txt/meta 三个后缀共用同一前缀)"""
for f in (entry.get("meta_file"), entry.get("html_file"), entry.get("txt_file")):
f = f or ""
if f.endswith(".meta.json"):
return f[:-10]
if f.endswith(".html"):
return f[:-5]
if f.endswith(".txt"):
return f[:-4]
return ""
def insert_results(run, results):
"""批量插入爬取结果 (增量); 文件只记基础名 base_file"""
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),
_base_of(r),
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,
base_file, image_count, image_files)
VALUES (%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)
# ---------------- 历史数据回填 (幂等) ----------------
def sync_all_history(task_ids=None):
"""把本地 JSON 中的历史任务/运行/结果全量回填数据库
可重复执行 (INSERT IGNORE + 唯一键去重)
task_ids: 指定只回填的任务ID列表, 默认全部
返回统计 dict
"""
import store as _store
stats = {"tasks": 0, "runs": 0, "results": 0}
for task in _store.load_tasks():
if task_ids and task["id"] not in task_ids:
continue
upsert_task(task)
stats["tasks"] += 1
for run in _store.get_runs(task["id"]):
upsert_run(run)
stats["runs"] += 1
results = run.get("results", [])
if results:
insert_results(run, results)
stats["results"] += len(results)
run["_db_count"] = len(results)
_store.save_run(task["id"], run) # 记录已同步标记, 避免运行中重复插入
return stats
+78 -19
View File
@@ -21,6 +21,7 @@ from playwright_stealth import Stealth
import notify
import store
import db
HERE = os.path.dirname(os.path.abspath(__file__))
DATA_DIR = os.path.join(HERE, "data")
@@ -37,13 +38,12 @@ _CHALLENGE_MARKS = [
def is_challenge_page(title, html):
"""判断是否仍在反爬验证页 (仅关键词启发, 避免误伤正常小页面)"""
low = html.lower()
t = (title or "").lower()
for mark in _CHALLENGE_MARKS:
if mark in t or mark in low:
return True
if len(html) < 5000 and ("<article" not in low and "<main" not in low):
return True
return False
@@ -77,6 +77,36 @@ def _settle_wait(page, timeout_s):
return True, page.title(), page.content()
_TRACKING_PARAMS = {
"utm_source", "utm_medium", "utm_campaign", "utm_term", "utm_content",
"fbclid", "gclid", "yclid", "mc_cid", "mc_eid", "ref", "ref_src",
}
def normalize_url(url):
"""URL 规范化 (用于去重): 去锚点/跟踪参数/尾部斜杠/默认端口, host 小写"""
try:
p = urllib.parse.urlparse(str(url))
host = (p.hostname or "").lower()
if not host:
return str(url)
port = ""
if p.port and p.port not in (80, 443):
port = f":{p.port}"
path = p.path or "/"
if len(path) > 1 and path.endswith("/"):
path = path.rstrip("/")
query = ""
if p.query:
kept = [kv for kv in p.query.split("&")
if kv.split("=", 1)[0].lower() not in _TRACKING_PARAMS]
if kept:
query = "?" + "&".join(kept)
return f"{p.scheme.lower()}://{host}{port}{path}{query}"
except Exception:
return str(url)
def filter_links(hrefs, seed_url, include=None, exclude=None,
same_domain=True, use_regex=False):
"""按规则过滤链接, 返回 (included, excluded); excluded 含排除原因"""
@@ -538,32 +568,48 @@ class CrawlJob:
return included
def _crawl_auto(self):
"""自动爬取: 持续递归 爬取->发现链接->爬取... 直到无新链接可爬或达到上限
- max_depth: 0=无限制, N=只爬 N 层
- max_pages: 0=无限制, N=安全上限
- 提取但未爬取的链接持久化到 auto.pending, 已爬集合持久化到 auto.visited
- 停止后再次运行从缓存队列继续爬; 起始网址每次运行都重新爬(不去重)
"""
run = self.run
auto = self.task.get("auto", {})
seed = auto.get("seed_url", "")
max_pages = int(auto.get("max_pages", 50) or 50)
max_depth = int(auto.get("max_depth", 2) or 2)
run["progress"]["total"] = max_pages
max_pages = int(auto.get("max_pages", 1000) or 0) # 0 = 无限制
max_depth = int(auto.get("max_depth", 0) or 0) # 0 = 无限制
out_dir = self._resolve_out_dir()
run["out_dir"] = out_dir
os.makedirs(out_dir, exist_ok=True)
# 恢复持久化状态: 已爬集合 + 上次未爬完的缓存队列
visited = set(auto.get("visited", []) or [])
pending = auto.get("pending", []) or []
# 起始网址每次运行都爬(不做去重), 缓存队列继续消费
queue = [(seed, 0, "")]
if pending:
queue.extend((p["url"], p.get("depth", 0), p.get("source", "")) for p in pending)
queued = set(visited)
for u, _d, _s in queue:
queued.add(normalize_url(u))
run["progress"]["total"] = len(queue)
self._persist()
p, browser, ctx, page, cookie_file = self._open_browser()
queue = [(seed, 0, "")] # (url, depth, 来源链接)
visited = set()
queued = set([seed])
idx = 0
try:
while queue and not self._stop.is_set():
self._wait_if_paused()
url, depth, src = queue.pop(0)
if url in visited:
continue
if len(visited) >= max_pages:
if max_pages > 0 and len(visited) >= max_pages:
self._log("info", f"达到最大页数上限 {max_pages}, 停止")
break
visited.add(url)
idx += 1
url, depth, src = queue.pop(0)
key = normalize_url(url)
if key in visited and url != seed: # 起始网址不去重, 其余已爬跳过
continue
visited.add(key)
idx = len(visited)
run["progress"]["current_url"] = url
run["progress"]["done"] = len(visited)
self._persist()
@@ -572,14 +618,27 @@ class CrawlJob:
run["results"].append(entry)
self._bump_stats(entry)
self._persist()
if entry["status"] == "OK" and depth < max_depth:
# 无深度限制或未达深度限制时持续发现链接
if entry["status"] == "OK" and (max_depth == 0 or depth < max_depth):
for link in self._discover_links(page):
if link not in visited and link not in queued:
queued.add(link)
lk = normalize_url(link)
if lk not in visited and lk not in queued:
queued.add(lk)
queue.append((link, depth + 1, url))
if entry["status"] == "OK":
self._delay()
run["progress"]["total"] = len(visited)
if not self._stop.is_set() and len(queue) == 0:
self._log("info", f"无新链接可爬, 任务结束 (共 {len(visited)} 页)")
finally:
self._close_browser(p, browser, ctx, cookie_file)
run["progress"]["total"] = len(visited)
# 无论完成/停止/异常, 都保存缓存队列与已爬集合, 便于下次继续
try:
auto["pending"] = [
{"url": u, "depth": d, "source": s} for u, d, s in queue]
auto["visited"] = list(visited)
store.upsert_task(self.task)
db.upsert_task(self.task)
except Exception:
pass
self._persist()
+1
View File
@@ -1,3 +1,4 @@
flask>=3.0
playwright>=1.40
playwright-stealth>=1.0
pymysql>=1.1
+5
View File
@@ -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
+49 -11
View File
@@ -79,10 +79,28 @@ function taskStatusBadge(t) {
return '<span class="badge">未运行</span>';
}
const GROUPS = [
{ key: "batch", label: "📄 批量爬取", desc: "一次性爬取指定网址列表" },
{ key: "scheduled", label: "⏰ 定时爬取", desc: "按间隔或 cron 表达式定时执行" },
{ key: "auto", label: "🤖 自动爬取", desc: "从起始网址自动发现链接,持续递归爬取到无新链接为止" },
];
function renderTasks() {
const box = $("taskList");
$("emptyState").classList.toggle("hidden", state.tasks.length > 0);
box.innerHTML = state.tasks.map(taskCard).join("");
box.innerHTML = GROUPS.map((g) => {
const items = state.tasks.filter((t) => t.mode === g.key);
if (!items.length) return "";
return `
<div class="group">
<div class="group-head">
<span class="group-title">${g.label}</span>
<span class="group-desc">${g.desc}</span>
<span class="group-count">${items.length} 个任务</span>
</div>
<div class="task-grid">${items.map(taskCard).join("")}</div>
</div>`;
}).join("");
}
function taskCard(t) {
@@ -174,10 +192,12 @@ async function readFileSmart(file) {
return text;
}
function importUrlsText(text) {
function importUrlsText(text, mode) {
const ta = $("taskForm").elements["urls"];
const clean = (s) => s.split(" #")[0].trim(); // 去掉行内注释 (URL 不含空格, 安全)
const existing = new Set(ta.value.split(/\r?\n/).map(clean).filter(Boolean));
const existing = new Set(
mode === "overwrite" ? [] : ta.value.split(/\r?\n/).map(clean).filter(Boolean)
);
const fresh = text.split(/\r?\n/).map(clean).filter(Boolean);
let added = 0;
for (const line of fresh) {
@@ -198,9 +218,14 @@ $("urlFileInput").onchange = async (e) => {
}
try {
const text = await readFileSmart(file);
const r = importUrlsText(text);
const mode = $("importMode").value;
const r = importUrlsText(text, mode);
formDirty = true;
toast(`已导入 ${file.name}:共 ${r.total} 行,新增 ${r.added} 条网址`);
if (mode === "overwrite") {
toast(`已导入 ${file.name}:共 ${r.total} 行(覆盖原列表,新增 ${r.added} 条)`);
} else {
toast(`已导入 ${file.name}:共 ${r.total} 行,新增 ${r.added} 条网址(追加)`);
}
} catch (err) {
toast("文件读取失败: " + err.message, true);
}
@@ -409,8 +434,8 @@ async function openEdit(tid) {
f.elements["exclude"].value = (t.auto.exclude || []).join("\n");
f.elements["same_domain"].checked = t.auto.same_domain !== false;
f.elements["use_regex"].checked = !!t.auto.use_regex;
f.elements["max_pages"].value = t.auto.max_pages ?? 50;
f.elements["max_depth"].value = t.auto.max_depth ?? 2;
f.elements["max_pages"].value = t.auto.max_pages ?? 1000;
f.elements["max_depth"].value = t.auto.max_depth ?? 0;
}
syncScheduleUI();
$("formHint").textContent = t.running
@@ -456,8 +481,8 @@ async function submitForm(e) {
exclude: splitLines(f.elements["exclude"].value),
same_domain: f.elements["same_domain"].checked,
use_regex: f.elements["use_regex"].checked,
max_pages: parseInt(f.elements["max_pages"].value) || 50,
max_depth: parseInt(f.elements["max_depth"].value) || 2,
max_pages: parseInt(f.elements["max_pages"].value) || 0,
max_depth: parseInt(f.elements["max_depth"].value) || 0,
};
}
if (mode === "scheduled") {
@@ -519,7 +544,10 @@ function renderDetail() {
<span>任务ID: <b>${t.id}</b></span>
<span>创建: <b>${fmtTime(t.created_at)}</b></span>
<span>输出: <b>${esc(cfg.out_dir || "out/" + t.id)}</b></span>
${t.mode === "auto" ? `<span>起始: <b>${esc((t.auto && t.auto.seed_url) || "")}</b></span>` : ""}
${t.mode === "auto" ? `
<span>起始: <b>${esc((t.auto && t.auto.seed_url) || "")}</b></span>
<span>待爬缓存: <b>${(t.auto && (t.auto.pending || []).length) || 0}</b> 条(停止后可继续爬取)</span>
<span>已爬: <b>${(t.auto && (t.auto.visited || []).length) || 0}</b> 页</span>` : ""}
</div>`;
let sched = "";
@@ -540,7 +568,8 @@ function renderDetail() {
<button class="btn" onclick="actTask('${t.id}','pause')">⏸ 暂停</button>
<button class="btn danger" onclick="actTask('${t.id}','stop')">⏹ 终止</button>`
: `<button class="btn primary" onclick="actTask('${t.id}','start')">▶ 立即执行</button>`}
${t.mode === "auto" ? `<button class="btn" onclick="probeTask('${t.id}')">🧪 试爬取</button>` : ""}
${t.mode === "auto" ? `<button class="btn" onclick="probeTask('${t.id}')">🧪 试爬取</button>
<button class="btn" ${t.running ? "disabled" : ""} onclick="clearCache('${t.id}')">🧹 清空缓存</button>` : ""}
<button class="btn" onclick="openEdit('${t.id}')">✏️ 编辑配置</button>
</div>`;
@@ -696,6 +725,15 @@ async function probeTask(tid) {
} catch (e) { toast(e.message, true); }
}
async function clearCache(tid) {
if (!confirm("清空待爬缓存与已爬记录?\n下次运行将从起始网址重新开始爬取。")) return;
try {
const r = await api(`/api/tasks/${tid}/clear_cache`, { method: "POST" });
toast(r.msg || "已清空");
openDetail(tid);
} catch (e) { toast(e.message, true); }
}
function renderProbe(r) {
const body = $("probeBody");
if (!r.ok) {
+10 -4
View File
@@ -77,7 +77,11 @@
<label>网址列表(每行一个,# 开头为注释)</label>
<div class="row2" style="align-items:center">
<button type="button" class="btn sm" id="btnImportUrls">📂 导入网址文件(.txt)</button>
<span class="card-line">支持 UTF-8 / GBK 编码,自动去重追加</span>
<select id="importMode" title="对已有网址的处理方式">
<option value="append" selected>追加(保留已有)</option>
<option value="overwrite">覆盖(清空已有)</option>
</select>
<span class="card-line">支持 UTF-8 / GBK,自动去重</span>
<input type="file" id="urlFileInput" accept=".txt,.csv,.urls,text/plain" hidden>
</div>
<textarea name="urls" rows="5" placeholder="https://www.example.com/&#10;https://www.example.com/page2"></textarea>
@@ -112,7 +116,7 @@
</div>
<div id="autoBox" class="hidden box">
<div class="field"><label>起始网址 *系统将自动发现符合规则的链接并爬取</label>
<div class="field"><label>起始网址 *(自动发现符合规则的链接,持续递归爬取直到无新链接可爬</label>
<div class="row2">
<input name="seed_url" placeholder="https://example.com/news" style="flex:1">
<button type="button" class="btn sm" id="btnProbe">🧪 试爬取</button>
@@ -129,8 +133,10 @@
<div class="field check"><label><input name="use_regex" type="checkbox"> 规则按正则匹配</label></div>
</div>
<div class="row2">
<div class="field"><label>最大页数</label><input name="max_pages" type="number" value="50"></div>
<div class="field"><label>最大爬取深度</label><input name="max_depth" type="number" value="2"></div>
<div class="field"><label>最大爬取深度(0=无限制,N=只爬 N 层)</label>
<input name="max_depth" type="number" min="0" value="0"></div>
<div class="field"><label>最大页数(0=无限制,默认 1000 作安全上限)</label>
<input name="max_pages" type="number" min="0" value="1000"></div>
</div>
</div>
</div>
+10
View File
@@ -85,6 +85,16 @@ main { padding: 20px 24px; max-width: 1500px; margin: 0 auto; }
.stat-num.num-run { color: var(--accent); }
.stat-label { font-size: 12px; color: var(--muted); margin-top: 3px; }
/* ---------- 分组展示 ---------- */
.group { display: flex; flex-direction: column; gap: 12px; margin-bottom: 24px; }
.group-head { display: flex; align-items: baseline; gap: 10px; flex-wrap: wrap; }
.group-title { font-size: 16px; font-weight: 700; }
.group-desc { font-size: 12px; color: var(--muted); }
.group-count {
font-size: 11px; padding: 2px 10px; border-radius: 10px;
background: var(--panel2); border: 1px solid var(--border); color: var(--muted);
}
/* ---------- task grid ---------- */
.task-grid { display: grid; grid-template-columns: repeat(auto-fill, minmax(380px, 1fr)); gap: 16px; }
.card {