在上游 MediaCrawler 之上新增一层: - 监控层 api/monitor/ —— 多博主/多笔记的定时采集、指标快照差分、报表、 企业微信通知。每轮采集写入独立目录,差分才成立。 - WebUI 登录鉴权 api/auth.py —— PBKDF2 口令 + 服务端会话,/api 全接口防护。 WebSocket 单独加依赖:BaseHTTPMiddleware 对 ws 作用域直接放行,覆盖不到。 - 全局平台切换 + 能力矩阵 —— 如实区分「爬虫模块支持」与「监控层已接线」, 未接通的平台直接拒绝建任务,而不是静默跑空。 - 监控库改用 MySQL 5.7(可回退 SQLite 供测试):逐表强制 utf8mb4 (服务端与库默认都是 latin1),启动校验所连 schema 以防写错库, 连接池 recycle + pre_ping 应对 MySQL 的 8 小时空闲断连。 修复上游缺陷: - xhs/core.py: 主页抓取失败会跳掉整个博主,导致一条作品都抓不到, 而那份资料只喂给一个空函数。改为尽力而为,失败不中断。 - xhs/login.py: cookie 登录只注入 web_session,冷启动签名会失败。 新增 INJECT_ALL_COOKIES 开关(默认关闭,原有行为不变)。 - requirements.txt: 补上 websockets。它在上游 pyproject.toml 里有声明、 这里漏了,导致 uvicorn 没有 WebSocket 能力,实时日志流从未工作。 改动过的上游文件清单及合并方式见 UPSTREAM.md。 测试:492 passed(另有 1 个既有的 Windows/gbk 上游测试失败,与本改动无关)
590 lines
19 KiB
Python
590 lines
19 KiB
Python
# -*- coding: utf-8 -*-
|
||
# Copyright (c) 2025 [email protected]
|
||
#
|
||
# This file is part of MediaCrawler project.
|
||
# Repository: https://github.com/NanmiCoder/MediaCrawler/blob/main/api/monitor/ingest.py
|
||
# GitHub: https://github.com/NanmiCoder
|
||
# Licensed under NON-COMMERCIAL LEARNING LICENSE 1.1
|
||
#
|
||
# 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则:
|
||
# 1. 不得用于任何商业用途。
|
||
# 2. 使用时应遵守目标平台的使用条款和robots.txt规则。
|
||
# 3. 不得进行大规模爬取或对平台造成运营干扰。
|
||
# 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。
|
||
# 5. 不得用于任何非法或不当的用途。
|
||
#
|
||
# 详细许可条款请参阅项目根目录下的LICENSE文件。
|
||
# 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。
|
||
|
||
"""Turn one run's crawled jsonl into snapshots and change events.
|
||
|
||
Pure-ish and offline testable: give it a directory of jsonl files, a run row and
|
||
a session, and it does the diffing. No network, no browser.
|
||
|
||
Correctness notes that drive the code below:
|
||
|
||
* Counts arrive as strings and may be abbreviated ("1.2万", "3亿"). A value that
|
||
cannot be parsed is stored as NULL, never 0 -- 0 would forge a large negative
|
||
delta on the next comparison.
|
||
* The comment endpoint has no time-sort, so only the platform's top-N window is
|
||
ever visible. A comment we have not seen before is therefore split into
|
||
"posted since last run" vs "seen for the first time", rather than claiming the
|
||
former always.
|
||
* A bad cookie does not make the crawler exit non-zero; it exits 0 having
|
||
fetched nothing. That is detected here as a suspected auth failure.
|
||
"""
|
||
|
||
import json
|
||
import re
|
||
from dataclasses import dataclass, field
|
||
from pathlib import Path
|
||
from typing import Any, Dict, List, Optional
|
||
|
||
from sqlalchemy import func, select
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
|
||
from tools.time_util import get_current_timestamp
|
||
|
||
from .platforms import PLATFORM_XHS
|
||
from .models import (
|
||
EVENT_AUTH_FAILURE,
|
||
EVENT_METRIC_DELTA,
|
||
EVENT_NEW_COMMENT_POSTED,
|
||
EVENT_NEW_COMMENT_SEEN,
|
||
EVENT_NEW_NOTE,
|
||
EVENT_NO_DATA,
|
||
EVENT_RUN_FAILED,
|
||
MonitorComment,
|
||
MonitorEvent,
|
||
MonitorNote,
|
||
MonitorNoteMetric,
|
||
MonitorRun,
|
||
MonitorTask,
|
||
RUN_FAILED,
|
||
RUN_PARTIAL,
|
||
RUN_SUCCESS,
|
||
)
|
||
|
||
_COUNT_UNITS = {
|
||
"": 1,
|
||
"万": 10_000,
|
||
"w": 10_000,
|
||
"W": 10_000,
|
||
"k": 1_000,
|
||
"K": 1_000,
|
||
"亿": 100_000_000,
|
||
}
|
||
_COUNT_RE = re.compile(r"^([\d.]+)\s*([万wWkK亿]?)$")
|
||
|
||
# Metric fields shared by the snapshot table and the delta comparison.
|
||
_METRIC_FIELDS = ("liked_count", "comment_count", "collected_count", "share_count")
|
||
|
||
|
||
def parse_count(value: Any) -> Optional[int]:
|
||
"""Parse an XHS interaction count into an int, or None if unintelligible.
|
||
|
||
Handles plain numbers, thousands separators, and the Chinese abbreviations
|
||
the platform actually returns ("1.2万" -> 12000, "3亿" -> 300000000).
|
||
"""
|
||
if value is None or isinstance(value, bool):
|
||
return None
|
||
if isinstance(value, int):
|
||
return value
|
||
if isinstance(value, float):
|
||
return int(value)
|
||
|
||
text = str(value).strip().replace(",", "").replace(" ", "")
|
||
if not text:
|
||
return None
|
||
|
||
match = _COUNT_RE.match(text)
|
||
if not match:
|
||
return None
|
||
|
||
try:
|
||
number = float(match.group(1))
|
||
except ValueError:
|
||
return None
|
||
|
||
return int(number * _COUNT_UNITS.get(match.group(2), 1))
|
||
|
||
|
||
# Windows reports hard process failures as NTSTATUS values, which surface in the
|
||
# UI as meaningless large integers (e.g. 3221225794 = 0xC0000142). Translating
|
||
# the ones we actually see saves the reader a hex-decoding detour.
|
||
_WINDOWS_EXIT_REASONS = {
|
||
0xC0000005: "进程访问冲突 (ACCESS_VIOLATION)",
|
||
0xC00000FD: "栈溢出 (STACK_OVERFLOW)",
|
||
0xC000013A: "进程被中断(控制台关闭或 Ctrl+C)",
|
||
0xC0000142: "进程初始化失败 (STATUS_DLL_INIT_FAILED),属启动环境异常,重启服务后重试",
|
||
0xC0000409: "栈缓冲区溢出 (STACK_BUFFER_OVERRUN)",
|
||
}
|
||
|
||
|
||
def describe_exit_code(code: int) -> str:
|
||
"""Render an exit code so a human can act on it."""
|
||
unsigned = code & 0xFFFFFFFF if code < 0 else code
|
||
reason = _WINDOWS_EXIT_REASONS.get(unsigned)
|
||
if reason:
|
||
return f"Crawler exited with code {code} (0x{unsigned:08X}): {reason}"
|
||
return f"Crawler exited with code {code}"
|
||
|
||
|
||
@dataclass
|
||
class IngestResult:
|
||
status: str
|
||
notes_fetched: int = 0
|
||
comments_fetched: int = 0
|
||
new_notes: int = 0
|
||
new_comments: int = 0
|
||
is_baseline: bool = False
|
||
error: Optional[str] = None
|
||
events: List[str] = field(default_factory=list)
|
||
|
||
|
||
def _read_jsonl(path: Path) -> List[Dict[str, Any]]:
|
||
"""Read a jsonl file, skipping blank or malformed lines."""
|
||
records: List[Dict[str, Any]] = []
|
||
if not path.exists():
|
||
return records
|
||
|
||
with path.open("r", encoding="utf-8") as handle:
|
||
for line in handle:
|
||
line = line.strip()
|
||
if not line:
|
||
continue
|
||
try:
|
||
item = json.loads(line)
|
||
except json.JSONDecodeError:
|
||
continue
|
||
if isinstance(item, dict):
|
||
records.append(item)
|
||
return records
|
||
|
||
|
||
def find_run_files(
|
||
out_dir: Path, platform: str = PLATFORM_XHS
|
||
) -> tuple[List[Path], List[Path]]:
|
||
"""Locate the contents/comments jsonl files a run produced.
|
||
|
||
The crawler writes ``{save_data_path}/{platform}/jsonl/{type}_{item}_{date}.jsonl``.
|
||
Glob rather than reconstructing the name: both the crawler type and the date
|
||
are runtime-dependent. Returns lists because a crawl crossing midnight
|
||
produces one file per day.
|
||
"""
|
||
jsonl_dir = out_dir / platform / "jsonl"
|
||
if not jsonl_dir.is_dir():
|
||
return [], []
|
||
|
||
return (
|
||
sorted(jsonl_dir.glob("*_contents_*.jsonl")),
|
||
sorted(jsonl_dir.glob("*_comments_*.jsonl")),
|
||
)
|
||
|
||
|
||
async def _emit(
|
||
session: AsyncSession,
|
||
run: MonitorRun,
|
||
event_type: str,
|
||
title: str,
|
||
*,
|
||
severity: str = "info",
|
||
target_kind: str = "",
|
||
target_id: str = "",
|
||
payload: Optional[Dict[str, Any]] = None,
|
||
) -> None:
|
||
session.add(
|
||
MonitorEvent(
|
||
task_id=run.task_id,
|
||
run_id=run.id,
|
||
type=event_type,
|
||
severity=severity,
|
||
target_kind=target_kind,
|
||
target_id=target_id,
|
||
title=title,
|
||
payload_json=json.dumps(payload or {}, ensure_ascii=False),
|
||
created_at=get_current_timestamp(),
|
||
)
|
||
)
|
||
|
||
|
||
async def _previous_run_started_at(
|
||
session: AsyncSession, task_id: int, run_id: int
|
||
) -> Optional[int]:
|
||
"""Started-at of the most recent earlier successful run, in ms."""
|
||
return await session.scalar(
|
||
select(MonitorRun.started_at)
|
||
.where(
|
||
MonitorRun.task_id == task_id,
|
||
MonitorRun.id != run_id,
|
||
MonitorRun.status.in_((RUN_SUCCESS, RUN_PARTIAL)),
|
||
MonitorRun.started_at.is_not(None),
|
||
# Same reasoning as _count_prior_successes: an empty run is a useless
|
||
# reference point for "was this comment posted since last time?".
|
||
MonitorRun.notes_fetched > 0,
|
||
)
|
||
.order_by(MonitorRun.id.desc())
|
||
.limit(1)
|
||
)
|
||
|
||
|
||
# How far back to look for proof that the stored login still works.
|
||
_AUTH_PROOF_WINDOW_MS = 6 * 60 * 60 * 1000
|
||
|
||
|
||
async def _another_task_succeeded_recently(session: AsyncSession, task_id: int) -> bool:
|
||
"""Whether a different task fetched data recently, proving the login is valid."""
|
||
since = get_current_timestamp() - _AUTH_PROOF_WINDOW_MS
|
||
count = await session.scalar(
|
||
select(func.count())
|
||
.select_from(MonitorRun)
|
||
.where(
|
||
MonitorRun.task_id != task_id,
|
||
MonitorRun.status == RUN_SUCCESS,
|
||
MonitorRun.started_at.is_not(None),
|
||
MonitorRun.started_at >= since,
|
||
)
|
||
)
|
||
return bool(count)
|
||
|
||
|
||
async def _count_prior_successes(session: AsyncSession, task_id: int, run_id: int) -> int:
|
||
return (
|
||
await session.scalar(
|
||
select(func.count())
|
||
.select_from(MonitorRun)
|
||
.where(
|
||
MonitorRun.task_id == task_id,
|
||
MonitorRun.id != run_id,
|
||
MonitorRun.status.in_((RUN_SUCCESS, RUN_PARTIAL)),
|
||
# A run that fetched nothing established no baseline. Without this
|
||
# check the first run that actually works after a failed one looks
|
||
# like a flood of newly discovered works.
|
||
MonitorRun.notes_fetched > 0,
|
||
)
|
||
)
|
||
) or 0
|
||
|
||
|
||
async def _ingest_notes(
|
||
session: AsyncSession,
|
||
run: MonitorRun,
|
||
records: List[Dict[str, Any]],
|
||
is_baseline: bool,
|
||
) -> int:
|
||
"""Upsert notes, write metric snapshots, and emit new-note/delta events."""
|
||
now = get_current_timestamp()
|
||
new_count = 0
|
||
|
||
for record in records:
|
||
note_id = record.get("note_id")
|
||
if not note_id:
|
||
continue
|
||
|
||
note = await session.scalar(
|
||
select(MonitorNote).where(
|
||
MonitorNote.task_id == run.task_id,
|
||
MonitorNote.note_id == note_id,
|
||
)
|
||
)
|
||
|
||
title = (record.get("title") or "")[:500]
|
||
raw_images = record.get("image_list") or ""
|
||
cover = raw_images.split(",")[0] if raw_images else ""
|
||
|
||
if note is None:
|
||
note = MonitorNote(
|
||
task_id=run.task_id,
|
||
note_id=note_id,
|
||
title=title,
|
||
note_url=record.get("note_url") or "",
|
||
cover=cover,
|
||
creator_hash=record.get("creator_hash") or "",
|
||
source_kind=record.get("type") or "",
|
||
published_at=_as_int(record.get("time")),
|
||
first_seen_run_id=run.id,
|
||
first_seen_at=now,
|
||
last_seen_run_id=run.id,
|
||
last_seen_at=now,
|
||
)
|
||
session.add(note)
|
||
new_count += 1
|
||
if not is_baseline:
|
||
await _emit(
|
||
session,
|
||
run,
|
||
EVENT_NEW_NOTE,
|
||
f"新作品:{title or note_id}",
|
||
target_kind="note",
|
||
target_id=note_id,
|
||
payload={"note_id": note_id, "title": title},
|
||
)
|
||
else:
|
||
# Only refresh descriptive fields; seen-tracking is updated below.
|
||
if title:
|
||
note.title = title
|
||
note.last_seen_run_id = run.id
|
||
note.last_seen_at = now
|
||
|
||
await _snapshot_metrics(session, run, note_id, record, now, is_baseline)
|
||
|
||
return new_count
|
||
|
||
|
||
def _as_int(value: Any) -> Optional[int]:
|
||
try:
|
||
return int(value)
|
||
except (TypeError, ValueError):
|
||
return None
|
||
|
||
|
||
async def _snapshot_metrics(
|
||
session: AsyncSession,
|
||
run: MonitorRun,
|
||
note_id: str,
|
||
record: Dict[str, Any],
|
||
now: int,
|
||
is_baseline: bool,
|
||
) -> None:
|
||
"""Write this run's metric snapshot and report any change vs the previous one."""
|
||
previous = await session.scalar(
|
||
select(MonitorNoteMetric)
|
||
.where(
|
||
MonitorNoteMetric.task_id == run.task_id,
|
||
MonitorNoteMetric.note_id == note_id,
|
||
MonitorNoteMetric.run_id != run.id,
|
||
)
|
||
.order_by(MonitorNoteMetric.run_id.desc())
|
||
.limit(1)
|
||
)
|
||
|
||
parsed = {name: parse_count(record.get(name)) for name in _METRIC_FIELDS}
|
||
|
||
session.add(
|
||
MonitorNoteMetric(
|
||
task_id=run.task_id,
|
||
note_id=note_id,
|
||
run_id=run.id,
|
||
captured_at=now,
|
||
liked_count=parsed["liked_count"],
|
||
comment_count=parsed["comment_count"],
|
||
collected_count=parsed["collected_count"],
|
||
share_count=parsed["share_count"],
|
||
raw_liked_count=str(record.get("liked_count") or ""),
|
||
raw_comment_count=str(record.get("comment_count") or ""),
|
||
raw_collected_count=str(record.get("collected_count") or ""),
|
||
raw_share_count=str(record.get("share_count") or ""),
|
||
)
|
||
)
|
||
|
||
if previous is None or is_baseline:
|
||
return
|
||
|
||
deltas = {}
|
||
for name in _METRIC_FIELDS:
|
||
old, new = getattr(previous, name), parsed[name]
|
||
# A None on either side means the value was unparseable; skip rather
|
||
# than report a bogus change.
|
||
if old is None or new is None or old == new:
|
||
continue
|
||
deltas[name] = {"from": old, "to": new, "delta": new - old}
|
||
|
||
if deltas:
|
||
summary = "、".join(
|
||
f"{_metric_label(name)} {info['from']}→{info['to']}"
|
||
for name, info in deltas.items()
|
||
)
|
||
await _emit(
|
||
session,
|
||
run,
|
||
EVENT_METRIC_DELTA,
|
||
f"互动数据变化:{summary}",
|
||
target_kind="note",
|
||
target_id=note_id,
|
||
payload={"note_id": note_id, "deltas": deltas},
|
||
)
|
||
|
||
|
||
def _metric_label(name: str) -> str:
|
||
return {
|
||
"liked_count": "点赞",
|
||
"comment_count": "评论",
|
||
"collected_count": "收藏",
|
||
"share_count": "分享",
|
||
}.get(name, name)
|
||
|
||
|
||
async def _ingest_comments(
|
||
session: AsyncSession,
|
||
run: MonitorRun,
|
||
records: List[Dict[str, Any]],
|
||
is_baseline: bool,
|
||
previous_run_started_at: Optional[int],
|
||
) -> int:
|
||
"""Upsert comments and emit events for ones never seen before."""
|
||
now = get_current_timestamp()
|
||
new_count = 0
|
||
|
||
for record in records:
|
||
comment_id = record.get("comment_id")
|
||
note_id = record.get("note_id")
|
||
if not comment_id or not note_id:
|
||
continue
|
||
|
||
exists = await session.scalar(
|
||
select(MonitorComment.id).where(
|
||
MonitorComment.task_id == run.task_id,
|
||
MonitorComment.note_id == note_id,
|
||
MonitorComment.comment_id == comment_id,
|
||
)
|
||
)
|
||
if exists is not None:
|
||
continue
|
||
|
||
create_time = _as_int(record.get("create_time"))
|
||
session.add(
|
||
MonitorComment(
|
||
task_id=run.task_id,
|
||
note_id=note_id,
|
||
comment_id=comment_id,
|
||
content=(record.get("content") or "")[:2000],
|
||
nickname=record.get("nickname") or "",
|
||
creator_hash=record.get("creator_hash") or "",
|
||
create_time=create_time,
|
||
like_count=parse_count(record.get("like_count")),
|
||
sub_comment_count=_as_int(record.get("sub_comment_count")) or 0,
|
||
parent_comment_id=record.get("parent_comment_id") or "",
|
||
first_seen_run_id=run.id,
|
||
first_seen_at=now,
|
||
)
|
||
)
|
||
new_count += 1
|
||
|
||
if is_baseline:
|
||
continue
|
||
|
||
# Without a time-sorted comment API we can only observe the top-N window,
|
||
# so distinguish a genuinely new comment from one that just surfaced.
|
||
posted = (
|
||
create_time is not None
|
||
and previous_run_started_at is not None
|
||
and create_time > previous_run_started_at
|
||
)
|
||
await _emit(
|
||
session,
|
||
run,
|
||
EVENT_NEW_COMMENT_POSTED if posted else EVENT_NEW_COMMENT_SEEN,
|
||
f"{'新评论' if posted else '新出现评论'}:{(record.get('content') or '')[:60]}",
|
||
target_kind="note",
|
||
target_id=note_id,
|
||
payload={
|
||
"note_id": note_id,
|
||
"comment_id": comment_id,
|
||
"create_time": create_time,
|
||
"nickname": record.get("nickname") or "",
|
||
},
|
||
)
|
||
|
||
return new_count
|
||
|
||
|
||
async def ingest_run(
|
||
session: AsyncSession,
|
||
run: MonitorRun,
|
||
task: MonitorTask,
|
||
out_dir: Path,
|
||
) -> IngestResult:
|
||
"""Ingest one finished run and return what changed.
|
||
|
||
Sets ``run.status``, ``run.is_baseline`` and the counters on the run row.
|
||
On a failed or untrustworthy run nothing is diffed -- the "seen" sets only
|
||
ever grow, so a partial run must never be allowed to look like deletions.
|
||
"""
|
||
# A non-zero exit is a genuine crash: trust nothing this run produced.
|
||
if run.exit_code not in (0, None):
|
||
run.status = RUN_FAILED
|
||
run.error_message = describe_exit_code(run.exit_code)
|
||
await _emit(
|
||
session,
|
||
run,
|
||
EVENT_RUN_FAILED,
|
||
f"采集进程异常退出(code={run.exit_code})",
|
||
severity="error",
|
||
payload={"exit_code": run.exit_code, "detail": run.error_message},
|
||
)
|
||
return IngestResult(status=RUN_FAILED, error=run.error_message)
|
||
|
||
contents_paths, comment_paths = find_run_files(out_dir, task.platform)
|
||
contents = [record for path in contents_paths for record in _read_jsonl(path)]
|
||
comments = [record for path in comment_paths for record in _read_jsonl(path)]
|
||
|
||
run.notes_fetched = len(contents)
|
||
run.comments_fetched = len(comments)
|
||
|
||
# A bad cookie does NOT fail the process: XHS cookie login is never validated,
|
||
# so an unauthenticated session just returns zero notes with exit 0 -- and
|
||
# usually does not even create an output file. Treating that as "the creator
|
||
# posted nothing" would silently hide login outages, which is exactly what
|
||
# monitoring exists to catch.
|
||
if not contents:
|
||
run.status = RUN_PARTIAL
|
||
|
||
# Blaming the cookie is only honest if nothing else is authenticating.
|
||
# A sibling task that just succeeded proves the login works, so the
|
||
# fault is with this target (bad/expired per-creator token, an empty
|
||
# account, or a page-structure change).
|
||
if await _another_task_succeeded_recently(session, run.task_id):
|
||
run.error_message = (
|
||
"Crawler produced no notes for this target, but other tasks "
|
||
"succeeded recently, so the login is probably fine"
|
||
)
|
||
await _emit(
|
||
session,
|
||
run,
|
||
EVENT_NO_DATA,
|
||
"本次未抓到任何作品:其他任务近期采集正常,登录态应该没问题,请检查该目标是否有效",
|
||
severity="warning",
|
||
payload={"out_dir": str(out_dir)},
|
||
)
|
||
else:
|
||
run.error_message = "Crawler produced no notes; the login cookie may have expired"
|
||
await _emit(
|
||
session,
|
||
run,
|
||
EVENT_AUTH_FAILURE,
|
||
"疑似登录态失效:本次未抓到任何作品,请检查 Cookie",
|
||
severity="error",
|
||
payload={"out_dir": str(out_dir)},
|
||
)
|
||
|
||
return IngestResult(
|
||
status=RUN_PARTIAL,
|
||
error=run.error_message,
|
||
comments_fetched=len(comments),
|
||
)
|
||
|
||
is_baseline = await _count_prior_successes(session, run.task_id, run.id) == 0
|
||
run.is_baseline = is_baseline
|
||
run.status = RUN_SUCCESS
|
||
run.error_message = None
|
||
|
||
previous_started_at = (
|
||
None if is_baseline else await _previous_run_started_at(session, run.task_id, run.id)
|
||
)
|
||
|
||
result = IngestResult(
|
||
status=RUN_SUCCESS,
|
||
notes_fetched=len(contents),
|
||
comments_fetched=len(comments),
|
||
is_baseline=is_baseline,
|
||
)
|
||
result.new_notes = await _ingest_notes(session, run, contents, is_baseline)
|
||
if task.enable_comments:
|
||
result.new_comments = await _ingest_comments(
|
||
session, run, comments, is_baseline, previous_started_at
|
||
)
|
||
|
||
run.new_notes = result.new_notes
|
||
run.new_comments = result.new_comments
|
||
return result
|