Files
MediaCrawler/api/monitor/ingest.py
T
butubb 06718a1351
Deploy VitePress site to Pages / build (push) Canceled after 0s
Deploy VitePress site to Pages / Deploy (push) Canceled after 0s
feat(monitor): 抖音接入博主监控
上游爬虫本身不缺抖音能力(三模式、四项指标、二级评论都与小红书对等、指标还是同名同列),
缺的全在监控层的适配。这次把「平台之间不一样」的管子集中到一个新模块,再把散落的
xhs 硬编码接上去。

* 新增 api/monitor/adapters.py:产物目录名、jsonl 字段别名、目标链接形态与正则、
  通知链接模板。不放进 platforms.py 是因为那个模块被 describe_all() 整个序列化进
  /api/config/platforms 交给前端,塞进正则和目录名会让爬虫内部细节漏进 API 载荷。
  代价是两个注册表可能漂移,用一条测试钉住「声明接通就必须有适配器」。
* 两个必须知道的坑,都在这版里处理掉了:
  1) 抖音的平台 id 是 dy,而 store 把产物写在 douyin/ 下(store/douyin/_store_impl.py:47)。
     不改就是 ingest 一个文件都读不到 —— 不报错,只是 0 条,然后被冒充成「疑似登录失效」。
  2) 抖音的作品没有 note_id(叫 aweme_id)、评论也用 aweme_id 指作品。ingest 第一步是
     `if not note_id: continue`,不映射就逐条全丢。
  另外抖音顶层评论的 parent_comment_id 是字符串 "0",归一成空串,免得前端多出悬空的父节点。
* 顺带把「东西抓到了、只是没落在期望目录里」单独识别出来。这类故障的现象和登录失效
  一模一样,按登录失效报会把人指去查完全错误的方向。
* 修两个既有 bug(今天只有小红书所以无害,加抖音就踩响):
  - service.py update_task 换目标时漏传 task.platform,回落到默认小红书
  - scheduler.py 取 cookie 没传 platform,抖音任务会读着小红书那份 cookie 不动
* 行为变更(已与用户确认):cookie 闸门改成「没 cookie 且没开 CDP」才跳过。
  CDP 模式下登录态来自被接管的浏览器,粘不粘 cookie 由不得它决定;不放行的话,
  选了「接管已有 Chrome」却没粘 cookie 的用户会看到任务永远不触发,而且不报错。
  副作用是开启了 CDP 的小红书任务也不再被该闸门拦住 —— 语义上是对的。
* 目标输入框的示例链接与措辞改由能力矩阵提供(notes_label 抖音说「作品」、小红书说
  「笔记」;「建议只填纯 ID」是小红书专属劝告,抖音链接不带令牌,不再显示)。

测试 +22 条(858 通过),其中最关键的是「抖音作品/评论不被静默丢弃」与「产物目录名
不等于平台 id」两条 —— 都是把最难查的失败模式钉死在回归网里。

注意:抖音这条路的**端到端尚未验证**,需要一份可用的抖音登录态(CDP 那台 Chrome 里
登录,或导出一份 cookie)。单测覆盖的是解析与入库,真实抓取还没跑过。
2026-10-10 14:55:28 +08:00

668 lines
23 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
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(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