Files
MediaCrawler/api/monitor/ingest.py
T
butubb e4affe9170
Deploy VitePress site to Pages / build (push) Canceled after 0s
Deploy VitePress site to Pages / Deploy (push) Canceled after 0s
feat(monitor): 运行历史要写清楚失败原因,不能只写「退出码 1」
用户的要求:运行历史的说明要写清楚。现在确实写不清楚 —— 抖音那次失败,运行历史里
只有一句 `Crawler exited with code 1`,而真正的报错 `DataFetchError: account blocked`
埋在子进程的 stderr 里,谁也看不到。

那两者本来是断开的两条路:子进程的输出只流向日志 WebSocket(前端 Terminal 看得到),
而监控层调 run_and_wait() 只拿得到一个退出码。

* crawler_manager 在 _push_log() 里留一份输出尾巴(80 行,每次 start 清空)——
  那是所有输出的唯一出口,挂这儿不会漏。新增 get_output_tail()。
* ingest 新增 diagnose_failure():倒着找第一行像异常的行(traceback 的末行),
  认不出就退回最后一行;管理器自己补的「Crawler exited with code」不是原因,排除掉。
* describe_exit_code() 接受这个原因并附在消息里;失败事件的标题也带上,这样企业微信
  通知和事件流不用翻日志就能看懂。
* runner 把尾巴交给 ingest;「超时/没起来」那条分支同样带上原因 —— -1 同时代表两种
  情况,而要查的东西完全不同。
* 运行历史那一格是截断的(240px),而失败原因现在有一整行 —— 补上 title 悬停显示,
  并放宽到 320px。没有悬停提示等于把最要紧的半句藏起来。

测试 +5:能挑出异常行、不会把管理器自己的话当成原因、没有输出时不报错、认不出时退回
最后一行、以及失败运行同时记下退出码与真因(含事件标题)。
2026-10-10 16:14:54 +08:00

727 lines
26 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, Sequence
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, cause: Optional[str] = None) -> str:
"""Render an exit code so a human can act on it.
只有退出码时信息量约等于零(`code 1` 什么都能是),所以把从输出尾巴里认出来的
异常一并附上 —— 运行历史里那一格显示的正是这句话。
"""
unsigned = code & 0xFFFFFFFF if code < 0 else code
reason = _WINDOWS_EXIT_REASONS.get(unsigned)
base = f"Crawler exited with code {code}"
if reason:
base = f"{base} (0x{unsigned:08X}): {reason}"
return f"{base};原因:{cause}" if cause else base
# 从爬虫输出里认出一行「异常」。Python 的 traceback 末行形如
# ``media_platform.douyin.exception.DataFetchError: account blocked``。
_EXCEPTION_LINE_RE = re.compile(r"^[\w.]*[A-Za-z](?:Error|Exception|Timeout)\b")
def diagnose_failure(output_tail: Optional[Sequence[str]]) -> Optional[str]:
"""从爬虫输出的末尾挑出最能说明问题的一行。
「退出码 1」等于什么都没说:真正的报错埋在子进程的 stderr 里。倒着找第一行看起来
像异常的行(traceback 的末行),找不到就退回最后一行有效输出。
"""
if not output_tail:
return None
lines = [line.strip() for line in output_tail if line and line.strip()]
# 管理器自己补的那两句不是爬虫的报错,别被当成失败原因。
noise = ("Crawler exited with code", "Crawler completed successfully")
lines = [line for line in lines if not line.startswith(noise)]
if not lines:
return None
for line in reversed(lines):
if _EXCEPTION_LINE_RE.match(line):
return line[:300]
return lines[-1][:300]
@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 "",
# 经 to_ms 换算:小红书给毫秒、抖音给秒,差 1000 倍。
published_at=adapter.to_ms(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
# 发布时间也刷。正常情况下它不会变,但**换算单位改过之后**(抖音是秒、
# 小红书是毫秒),已经入库的那批只能靠重采修回来。
published = adapter.to_ms(adapter.note_field(record, "published_at"))
if published is not None:
note.published_at = published
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
# 时间同理:单位换算修好之后,老数据要重采才能纠正。
created = adapter.to_ms(adapter.comment_field(record, "create_time"))
if created is not None:
existing.create_time = created
continue
create_time = adapter.to_ms(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,
output_tail: Optional[Sequence[str]] = None,
) -> 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
cause = diagnose_failure(output_tail)
run.error_message = describe_exit_code(run.exit_code, cause)
title = f"采集进程异常退出(code={run.exit_code})"
if cause:
# 标题里也带上真因:企业微信通知和事件流都只看这一行,不写就还得去翻日志。
title = f"{title}:{cause}"
await _emit(
session,
run,
EVENT_RUN_FAILED,
title,
severity="error",
payload={
"exit_code": run.exit_code,
"detail": run.error_message,
"cause": cause,
},
)
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