统一异构数据源的可靠路径是”一套标准数据模型 + 每个源一个适配器”:先定义与任何平台无关的标准记录结构,再为 Twitter、B 站等每个源写一个薄适配器把原始响应翻译成它,下游的分析、去重、推送只认标准模型。AI 热点监控工具用这一结构替换了早期堆在业务层的条件分支,新增一个平台只需增加一个适配器文件,检索与推送链路不动。下面给出模型定义、适配器实现与归一化的具体规则。
一、异构数据源的差异到底在哪
把几个平台的搜索响应摆在一起,差异集中在五个层面。
| 差异维度 | Twitter 类接口 | B 站类接口 | 归一化处理 |
|---|---|---|---|
| 字段命名 | 推文文本、转推数、点赞数 | 视频标题、播放量、弹幕数 | 映射到标准模型同名字段 |
| 时间格式 | 带时区的字符串时间戳 | 秒级 Unix 时间戳 | 统一转 UTC 的 ISO 8601 |
| 热度指标 | 互动量(转推 + 点赞) | 播放量、弹幕量 | 折算为统一热度分 |
| 主键形式 | 平台自有长整型 ID | 稿件号与视频 ID 并存 | 源名 + 源内 ID 做确定性哈希 |
| 分页方式 | 游标翻页 | 页码加页长 | 适配器内部消化,对外只暴露批次 |
差异不只在字段名。同一条”热点”在不同平台的语义粒度也不同:一边是单条短文本,一边是带时长与分区的视频。标准模型要能无损承载判断热度所需的信息,又不能为迁就某个平台而膨胀。
二、标准数据模型的字段设计
标准模型是整个系统的契约。它定义一次,所有下游依赖它,不依赖任何源的原始结构。
核心字段固定存在:全局唯一 ID、源名、源内 ID、标题、正文、作者匿名引用、发布时间、采集时间、原文链接、内容类型。可选字段按源填充:热度原值、互动细项、媒体时长、话题标签。
两个字段值得单独说明。
原始载荷副本要完整保留。平台调整响应结构时,适配器可能静默失效,这份未加工的记录是排查依据,不要在入库时丢弃。
字段覆盖度是可观测性抓手。适配器在转换末尾计算”可选字段实际填充的比例”,把数据缺失显式化,而不是让空值散落在下游各处的判空逻辑里。覆盖度按源按天聚合后,一旦某个源从九成掉到两成,说明接口结构变了,能在用户察觉前触发告警。
from datetime import datetime
from enum import Enum
from typing import Any, Optional
from pydantic import BaseModel, Field, HttpUrl
class ContentType(str, Enum):
TEXT = "text"
VIDEO = "video"
IMAGE = "image"
class HotItem(BaseModel):
"""跨源统一的热点标准记录,下游只消费这个结构。"""
uid: str # 源名 + 源内 ID 的确定性哈希
source: str # twitter / bilibili / ...
source_id: str
title: str
content: str = ""
author_ref: str = "" # 匿名指纹,不含可识别标识
content_type: ContentType = ContentType.TEXT
published_at: datetime # 统一 UTC
collected_at: datetime
url: Optional[HttpUrl] = None
heat_raw: int = 0 # 源内原始热度值
heat_score: int = 0 # 归一化后的 0-100 分
tags: list[str] = Field(default_factory=list)
coverage: float = 0.0 # 可选字段填充比例
raw: dict[str, Any] = Field(default_factory=dict) # 原始载荷留档
三、每个源一个适配器
适配器是纯函数式的转换:单个源的原始输入进,一条标准记录出。它承担这个源的全部特殊性,不承担别的。
字段改名、类型转换、枚举归一、覆盖度计算都在适配器内完成。平台改了字段名,只改这一个文件;分析与推送模块不知道发生过变化。
import hashlib
from abc import ABC, abstractmethod
from datetime import datetime, timezone
REGISTRY: dict[str, "SourceAdapter"] = {}
def make_uid(source: str, source_id: str) -> str:
"""确定性 ID:同一条数据在多次采集中保持一致,便于幂等写入。"""
return hashlib.sha1(f"{source}:{source_id}".encode()).hexdigest()
class SourceAdapter(ABC):
source: str
optional_fields = ("content", "url", "tags", "heat_raw")
@abstractmethod
def to_item(self, raw: dict) -> HotItem: ...
def coverage_of(self, item: HotItem) -> float:
filled = sum(1 for f in self.optional_fields if getattr(item, f))
return round(filled / len(self.optional_fields), 2)
def __init_subclass__(cls, **kw):
super().__init_subclass__(**kw)
REGISTRY[cls.source] = cls()
class BilibiliAdapter(SourceAdapter):
source = "bilibili"
def to_item(self, raw: dict) -> HotItem:
item = HotItem(
uid=make_uid(self.source, str(raw["bvid"])),
source=self.source,
source_id=str(raw["bvid"]),
title=raw.get("title", "").strip(),
content=raw.get("desc", "").strip(),
content_type=ContentType.VIDEO,
# 秒级时间戳统一为带时区的 UTC 时间
published_at=datetime.fromtimestamp(raw["pubdate"], tz=timezone.utc),
collected_at=datetime.now(timezone.utc),
url=f"https://www.bilibili.com/video/{raw['bvid']}",
heat_raw=int(raw.get("stat", {}).get("view", 0)),
tags=[t for t in (raw.get("tag") or "").split(",") if t],
raw=raw,
)
item.coverage = self.coverage_of(item)
return item
def normalize(source: str, batch: list[dict]) -> list[HotItem]:
adapter = REGISTRY[source]
out = []
for raw in batch:
try:
out.append(adapter.to_item(raw))
except Exception:
# 单条失败不拖垮整批,转换异常单独记录待排查
continue
return out
四、接入一个新数据源的步骤
流程固定,可复用:
- 抓取该源的真实响应样本,逐字段对照标准模型列出映射关系;
- 确认时间字段的单位与时区,统一转成 UTC 的 ISO 表示;
- 确定源内主键,用源名加主键生成确定性 ID,保证重复采集幂等;
- 编写适配器子类实现转换,把可选字段缺失情况纳入覆盖度计算;
- 用样本数据跑校验,标准模型的类型约束会拦下不合法字段;
- 在采集调度里注册该源的轮询周期与限流配额;
- 上线后观察覆盖度曲线,确认适配器在真实流量下没有静默失效。
五、归一化的三条硬规则
5.1 枚举收敛到小集合
内容类型、平台名、状态这类字段必须落到一个可比的小集合里。若放任各源自带的表述并存,下游每处判断都要写一串等价映射,规则会随源的增加而膨胀。
5.2 时间与数值单位单一
时间只用 UTC 的 ISO 字符串对外传递,展示时区交前端处理。热度数值先记原值,再按平台内分位数折算成 0–100 的分数;跨源排序只用分数,不用原值。
5.3 跨源去重放在标准模型之后
同一个事件常在多平台同时出现。去重逻辑基于标准模型的标题指纹与链接归一化结果,运行在适配器下游。这样新增源不需要改动去重规则,只要它输出的是标准记录。
六、校验与异常兜底
标准模型带类型约束,转换结果在构造时就被校验,脏数据不会流入存储。单条转换失败只跳过该条并留档,不中断整批采集。原始载荷副本配合覆盖度指标,构成排查接口变更的两个抓手:覆盖度告警提示”有东西坏了”,原始载荷回答”具体坏在哪个字段”。
常见问题(FAQ)
Q1:标准模型字段不够用怎么办?
新增可选字段并给默认值,旧适配器无需改动,覆盖度会自然反映填充情况。
Q2:某个源的接口结构变了如何发现?
监控该源字段覆盖度曲线,出现明显下跌即触发告警,再查原始载荷定位。
Q3:适配器要不要处理去重?
不处理。去重与打分运行在标准模型之上,适配器只负责单源到标准记录的翻译。