两件都是「作品栏里把这东西认出来」的延伸: * **账号级指标**:作品列表只会说「这条涨了多少赞」,说不了「这个人整个 账号的粉丝在涨还是在掉」。抖音的资料接口本来就有粉丝数/总获赞/作品数, 每轮顺手记一条快照(`monitor_creator_stat`,粒度 = 任务×博主×轮次, 和作品指标同形)。组头显示最近一条。 快照在「一条作品都没采到」的早退**之前**落:作品列表被风控挡住的那一轮, 正是「粉丝还在涨、但新作品没在发现」最该被看见的时刻。 * **作品备注**:博主备注回答「这个账号是谁」,这条回答「这条我要盯着」。 一个博主底下常常只有一两件值得盯的作品,所以不能合并成一条。键取 (platform, note_id),跨任务共用一份。 两边都守住同一条口径:**不知道就是 null,不写成 0** —— 0 在趋势图上是一条 砸到底的线,和「还没采到」是两回事。
801 lines
29 KiB
Python
801 lines
29 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, 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,
|
||
MonitorCreatorStat,
|
||
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 find_profile_files(out_dir: Path, platform: str = PLATFORM_XHS) -> List[Path]:
|
||
"""博主**账号级**指标那几行 jsonl(``creator_profile_*.jsonl``)。
|
||
|
||
单独一个函数而不是往 ``find_run_files`` 的返回值里塞第三个列表:那个返回值的两个
|
||
位置是有意义的(contents/comments),加一个会把所有调用点和解包语句都牵动一遍,
|
||
而这份产物是**可选**的 —— 小红书那条路(爬虫进程)根本不产生它。
|
||
"""
|
||
jsonl_dir = out_dir / adapters.artifact_dir(platform) / "jsonl"
|
||
if not jsonl_dir.is_dir():
|
||
return []
|
||
return sorted(jsonl_dir.glob("*_profile_*.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
|
||
# 同一轮里重复出现的作品只处理一次。**指标快照的唯一键是 (task_id, note_id, run_id)**,
|
||
# 同一件作品在一轮里进来两次会让第二次插入直接撞键、整个 run 崩掉 —— 产物里重复并不
|
||
# 罕见(多个目标指向同一个人、或退化路径重复刷新)。
|
||
seen_in_run: set = set()
|
||
|
||
for record in records:
|
||
note_id = adapter.note_field(record, "note_id")
|
||
if not note_id:
|
||
continue
|
||
if note_id in seen_in_run:
|
||
continue
|
||
seen_in_run.add(note_id)
|
||
|
||
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_creator_stats(
|
||
session: AsyncSession,
|
||
run: MonitorRun,
|
||
records: Sequence[Dict[str, Any]],
|
||
) -> int:
|
||
"""把这一轮问到的博主账号级指标落成快照,返回条数。
|
||
|
||
和作品指标一样是**每轮一条**:账号级的粉丝数是缓慢变化的量,「今天比昨天多了 300」
|
||
才是有用的信号,单看一个绝对值没有意义 —— 所以这里只管记,分析交给查询端。
|
||
|
||
**没解析出来的值留 NULL,不写 0**:0 在趋势图上是一条砸到底的线,和「不知道」完全是
|
||
两回事(见 ``parse_count`` 的注释)。
|
||
"""
|
||
now = get_current_timestamp()
|
||
written = 0
|
||
seen: set = set()
|
||
|
||
for record in records:
|
||
creator_hash = str(record.get("creator_hash") or "").strip()
|
||
if not creator_hash or creator_hash in seen:
|
||
# 一个任务可以配多个目标,退化路径下它们可能指向同一个博主 —— 而唯一键是
|
||
# (task_id, creator_hash, run_id),重复插入会撞键把整轮炸掉。
|
||
continue
|
||
seen.add(creator_hash)
|
||
|
||
session.add(
|
||
MonitorCreatorStat(
|
||
task_id=run.task_id,
|
||
run_id=run.id,
|
||
creator_hash=creator_hash,
|
||
nickname=str(record.get("nickname") or "")[:128],
|
||
fans=parse_count(record.get("fans")),
|
||
total_favorited=parse_count(record.get("total_favorited")),
|
||
works_count=parse_count(record.get("works")),
|
||
following=parse_count(record.get("following")),
|
||
captured_at=now,
|
||
)
|
||
)
|
||
written += 1
|
||
|
||
return written
|
||
|
||
|
||
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)
|
||
|
||
# 账号级快照**在「一条作品都没采到」的早退之前**落。博主的粉丝数并不会因为他这个
|
||
# 月的新作品列表被风控挡住就不存在 —— 那正是最该看到「粉丝还在涨、但新作品没在发现」
|
||
# 的时刻,跳过它等于在最需要它的那轮把数据丢掉。
|
||
profiles = [
|
||
record
|
||
for path in find_profile_files(out_dir, task.platform)
|
||
for record in _read_jsonl(path)
|
||
]
|
||
await _ingest_creator_stats(session, run, profiles)
|
||
|
||
# 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
|