Files
MediaCrawler/api/monitor/ingest.py
T
butubb cd85587f00
Deploy VitePress site to Pages / build (push) Canceled after 0s
Deploy VitePress site to Pages / Deploy (push) Canceled after 0s
feat(privacy): 关掉昵称脱敏 —— 脱敏有损,撞名就分不出博主
需求:评论栏/作品栏要能分清是哪个博主。修好分组字段之后名字仍带星号,因为上游作为
教学版默认对昵称做中间脱敏(首尾各留 1 字,中间星号)。这个脱敏是**有损**的:

  「张三」和「张四」都变成「张*」
  「小明老师」和「小刚老师」都变成「小***师」

而本仓库的用途是监控一批公开创作者账号,分清谁是谁正是这一层要干的事。所以关掉它。

* config/base_config.py 新增 MASK_NICKNAME = False(和 INJECT_ALL_COOKIES 一样是个
  开关,不是删代码 —— 改回 True 就恢复上游行为)。
* tools/user_hash.py 的 mask_nickname 读这个开关,关闭时原样返回。读的是模块属性而
  不是导入值,测试才能 monkeypatch。脱敏实现本身一字未动。
* 顺带修一个数据陈旧问题:评论是去重后直接 continue 的,昵称只在首次入库时写一次,
  于是开关一改(或评论者改名)老评论永远停在旧值 —— 而重采是唯一能拿到新值的途径。
  现在已存在的评论会跟着刷新昵称(作品那边的 creator_name 早就是这么做的)。
* anonymous 的 creator_hash 保持不变:那是分组用的稳定键,不是显示名。

测试:
* 三个隐私套件 + weibo 的 autouse fixture 强制把开关打开 —— 它们验的是**脱敏机制
  本身**,机制仍然必须正确,所以显式打开来测,而不是让它们随部署配置漂。
* test_mask_and_hash_tools 改成两个方向都覆盖(开着脱敏 / 关着脱敏)。
* test_tieba_extractor.py 里 8 处字面量的脱敏期望值换成真实昵称 —— 提取器现在就是
  返回原文的,期望值理应跟着改(这一条是行为变更的直接后果,不是测试放宽)。
* 新增一条:已入库的评论昵称会随重采刷新(且不会因刷新而重复插入)。
2026-10-10 15:48:15 +08:00

674 lines
24 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# -*- 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 . import adapters
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.
``platform`` 是**监控层的平台 id**,而爬虫落盘的目录名未必同名(抖音的 id 是
``dy``、目录是 ``douyin``),所以这里经 adapters 解析 —— 调用方不必知道这个差异。
"""
jsonl_dir = out_dir / adapters.artifact_dir(platform) / "jsonl"
if not jsonl_dir.is_dir():
return [], []
return (
sorted(jsonl_dir.glob("*_contents_*.jsonl")),
sorted(jsonl_dir.glob("*_comments_*.jsonl")),
)
def _misplaced_output_dirs(out_dir: Path, expected: str) -> List[str]:
"""在 out_dir 下找「有产物、但目录名不是期望的那个」的目录。
这是专门为**最难查的那类故障**准备的:产物目录名与平台对不上时,ingest 一个文件
都找不到,现象和「登录态失效」一模一样 —— 而实际上登录好好的、数据也抓到了,
只是没人去对的地方读。上游哪天改了 store 的目录名,这里能直接把实情说出来。
"""
found = []
try:
children = list(out_dir.iterdir())
except OSError:
return found
for child in children:
if child.name == expected or not child.is_dir():
continue
try:
if any(child.glob("jsonl/*_contents_*.jsonl")):
found.append(child.name)
except OSError:
continue
return sorted(found)
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,
adapter: adapters.PlatformAdapter,
) -> int:
"""Upsert notes, write metric snapshots, and emit new-note/delta events.
记录里的字段一律经 ``adapter`` 读。抖音的作品没有 ``note_id``(叫 ``aweme_id``),
按名字硬取的话每条记录都会在下面第一行被 continue 掉 —— 一条都不报错地全丢。
"""
now = get_current_timestamp()
new_count = 0
for record in records:
note_id = adapter.note_field(record, "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 = (adapter.note_field(record, "title") or "")[:500]
cover = adapter.cover(record)
if note is None:
note = MonitorNote(
task_id=run.task_id,
note_id=note_id,
title=title,
note_url=adapter.note_field(record, "note_url") or "",
cover=cover,
creator_hash=adapter.note_field(record, "creator_hash") or "",
creator_name=adapter.note_field(record, "creator_name") or "",
source_kind=adapter.note_field(record, "source_kind") or "",
published_at=_as_int(adapter.note_field(record, "published_at")),
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
# 封面地址**带签名、会过期**,所以每轮都用最新的覆盖它。原先只在首次入库
# 时写一次,结果旧作品的封面地址烂在库里 —— 隔天开始全是 403,而且再怎么
# 重跑也修不回来。落盘那份由 covers.cache_pending 负责(网络操作不在本模块)。
if cover:
note.cover = cover
# 昵称也要刷:作者改昵称是常事,只在首次入库写一次会一直显示旧的。
creator_name = adapter.note_field(record, "creator_name")
if creator_name:
note.creator_name = creator_name
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],
adapter: adapters.PlatformAdapter,
) -> int:
"""Upsert comments and emit events for ones never seen before.
与作品同理,评论记录也要经 ``adapter`` 读:抖音的评论用 ``aweme_id`` 指作品。
"""
now = get_current_timestamp()
new_count = 0
for record in records:
comment_id = adapter.comment_field(record, "comment_id")
note_id = adapter.comment_field(record, "note_id")
if not comment_id or not note_id:
continue
existing = await session.scalar(
select(MonitorComment).where(
MonitorComment.task_id == run.task_id,
MonitorComment.note_id == note_id,
MonitorComment.comment_id == comment_id,
)
)
if existing is not None:
# 昵称要跟着刷,不能只写一次。评论是去重后直接 continue 的,若不刷新,
# 脱敏开关一改(或评论者改了昵称),已经入库的老评论会永远停在旧值上 ——
# 而重采是唯一能拿到新值的途径。作品那边的 creator_name 同理。
refreshed = adapter.comment_field(record, "creator_name")
if refreshed:
existing.nickname = refreshed
continue
create_time = _as_int(adapter.comment_field(record, "create_time"))
session.add(
MonitorComment(
task_id=run.task_id,
note_id=note_id,
comment_id=comment_id,
content=(adapter.comment_field(record, "content") or "")[:2000],
nickname=adapter.comment_field(record, "creator_name") or "",
creator_hash=adapter.comment_field(record, "creator_hash") or "",
create_time=create_time,
like_count=parse_count(adapter.comment_field(record, "like_count")),
sub_comment_count=_as_int(
adapter.comment_field(record, "sub_comment_count")
)
or 0,
parent_comment_id=adapter.parent_comment_id(record),
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)
adapter = adapters.adapter(task.platform)
subdir = adapters.artifact_dir(task.platform)
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
# 先排除「东西抓到了,只是没落在我们找的那个目录里」。这种故障的现象和登录失效
# 一模一样,但登录其实是好的 —— 按登录失效报会把人指到完全错的方向去查。
misplaced = _misplaced_output_dirs(out_dir, subdir)
if misplaced:
run.error_message = (
f"crawler wrote into {misplaced} but platform {task.platform} "
f"expects {subdir}"
)
await _emit(
session,
run,
EVENT_NO_DATA,
f"采集产物目录与平台不匹配(实际 {misplaced}、期望 {subdir}),本次未读到任何作品",
severity="error",
payload={
"out_dir": str(out_dir),
"expected": subdir,
"found": misplaced,
},
)
return IngestResult(
status=RUN_PARTIAL,
error=run.error_message,
comments_fetched=len(comments),
)
# 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, adapter)
if task.enable_comments:
result.new_comments = await _ingest_comments(
session, run, comments, is_baseline, previous_started_at, adapter
)
run.new_notes = result.new_notes
run.new_comments = result.new_comments
return result