7.11 自建系统架构:从 API 到你的数据中台
架构这一篇讲怎么把前面所有零件组装成一个可持续运行的系统。
一、五层架构
| 层级 | 职责 | 关键设计 |
|---|---|---|
| ① 采集层 | 调用 API 获取原始数据 | 限流、重试、去重、缓存 |
| ② 存储层 | 持久化原始数据与衍生指标 | 分表、索引、避免重复拉取 |
| ③ 计算层 | 构造高阶指标 | 可重算、可版本化 |
| ④ 应用层 | 生成分析、内容、告警 | 解耦,可独立替换 |
| ⑤ 调度层 | 定时触发全流程 | 幂等、可重试、有日志 |
📊 架构图占位
图 7.3:自建分析系统五层架构
此处放置五层架构图,展示采集→存储→计算→应用→调度的数据趋势分布,并标注每层的技术选型建议。
📐 建议尺寸:16:9 高清矢量图🔴 重点标注:存储层与计算层之间的「原始数据 / 衍生指标」分离
二、采集层设计
三个必须处理的工程问题
import time, random, requests, logging
from functools import wraps
def retry(max_attempts=3, base_delay=1.0):
"""带实时指数退避的重试装饰器"""
def deco(fn):
@wraps(fn)
def wrapper(*args, **kwargs):
for attempt in range(1, max_attempts + 1):
try:
return fn(*args, **kwargs)
except requests.HTTPError as e:
# 4xx 类错误不重试(除了 429)
code = getattr(e.response, 'status_code', None)
if code and 400 <= code < 500 and code != 429:
raise
if attempt == max_attempts:
raise
# 实时指数退避 + 随机抖动,避免惊群
delay = base_delay * (2 ** (attempt - 1)) + random.random()
logging.warning(f'第 {attempt} 次失败,{delay:.1f}s 后重试: {e}')
time.sleep(delay)
return None
return wrapper
return deco
@retry(max_attempts=3)
def fetch_json(url, headers, params=None, timeout=20):
r = requests.get(url, headers=headers, params=params, timeout=timeout)
r.raise_for_status()
return r.json()
限流保护
class RateLimiter:
"""简易令牌桶限流器:控制调用频率,避免触发 429"""
def __init__(self, min_interval=1.2):
self.min_interval = min_interval
self._last = 0.0
def wait(self):
elapsed = time.time() - self._last
if elapsed < self.min_interval:
time.sleep(self.min_interval - elapsed)
self._last = time.time()
三、存储层设计
分层存储原则
| 数据 | 存储方式 | 原因 |
|---|---|---|
| 原始响应 | 原始 JSON 落盘 / 对象存储 | 可回溯、可重新解析 |
| 结构化数据 | SQLite / PostgreSQL 表 | 便于查询与关联 |
| 衍生指标 | 单独的表,带计算版本号 | 算法迭代时可重算 |
为什么衍生指标要带版本号
你的算法一定会改。
如果衍生指标没有版本标记,改动算法后你无法区分「哪一批数字是旧算法算的」。
建议在指标表里加一列calc_version,每次算法变更就递增。
这样你可以:重算时只处理旧版本数据、对比新旧算法效果、必要时回滚。
import sqlite3
def init_db(path='system.db'):
"""初始化数据表"""
conn = sqlite3.connect(path)
conn.executescript('''
CREATE TABLE IF NOT EXISTS raw_responses (
id INTEGER PRIMARY KEY AUTOINCREMENT,
endpoint TEXT NOT NULL,
cache_key TEXT UNIQUE,
payload TEXT NOT NULL,
fetched_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS metrics (
id INTEGER PRIMARY KEY AUTOINCREMENT,
fixture_id INTEGER NOT NULL,
team_id INTEGER NOT NULL,
metric_name TEXT NOT NULL,
metric_value REAL,
calc_version TEXT NOT NULL,
computed_at TEXT NOT NULL,
UNIQUE(fixture_id, team_id, metric_name, calc_version)
);
CREATE INDEX IF NOT EXISTS idx_metrics_lookup
ON metrics(fixture_id, metric_name);
''')
return conn
四、计算层设计
核心原则:可重算
class MetricCalculator:
"""指标计算器基类
每个指标一个类,便于单独测试与版本管理。
"""
version = 'v1'
def __init__(self, conn):
self.conn = conn
def compute(self, fixture_id, team_id, context):
"""返回本指标的值,子类实现"""
raise NotImplementedError
def save(self, fixture_id, team_id, name, value):
self.conn.execute(
'INSERT OR REPLACE INTO metrics '
'(fixture_id, team_id, metric_name, metric_value, calc_version, computed_at) '
'VALUES (?, ?, ?, ?, ?, datetime("now"))',
(fixture_id, team_id, name, value, self.version)
)
class AttackStrength(MetricCalculator):
"""进攻强度指标(教学示例)"""
version = 'v1'
def compute(self, fixture_id, team_id, context):
recent = context.get('recent_matches', [])
if len(recent) < 5:
return None # 样本不足,明确返回 None 而不是猜
decayed = decayed_stats(recent, decay=0.9)
return round(decayed['xg_avg'], 4)
五、调度层设计
幂等性
定时任务必须可以重复执行而不产生副作用:
def daily_job(fixture_ids, conn):
"""每日任务:幂等,重复执行不会产生重复数据"""
for fid in fixture_ids:
# 先检查是否已处理过
cur = conn.execute(
'SELECT COUNT(*) FROM metrics WHERE fixture_id = ? AND calc_version = ?',
(fid, AttackStrength.version)
)
if cur.fetchone()[0] > 0:
logging.info(f'fixture {fid} 已处理,跳过')
continue
try:
run_pipeline(fid, conn)
except Exception as e:
# 单个失败不影响整体
logging.error(f'fixture {fid} 处理失败: {e}')
continue
可观测性
| 需要记录 | 用途 |
|---|---|
| 每次调用的接口与耗时 | 性能与成本分析 |
| 点数消耗累计 | 成本预警 |
| 计算失败与原因 | 排障 |
| 指标样本量 | 判断置信度 |
六、应用层:解耦
把「算指标」和「产出内容」彻底分开。
好处:
· 换 AI 模型不影响算指标的逻辑;
· 指标可以用于多种输出(文章、告警、图表);
· 出问题容易定位到是哪一层。
七、技术选型建议
| 规模 | 推荐方案 |
|---|---|
| 个人实验 | Python + SQLite + 文件缓存 + 手动运行 |
| 小规模生产 | Python + SQLite + APScheduler + 日志文件 |
| 多人使用 | Python + PostgreSQL + Celery + 监控告警 |
不要一开始就上重型架构。很多需求用 SQLite + 一个 Python 脚本就够了。先跑起来,再按需要升级。
下一步该看什么
- 最小可用示例 → 7.12 30 行代码预警机器人
- 全自动化 → 7.13 自动化流水线
足球赛事前瞻 | AI 智能分析 | 足球数据解读 - 球小策