2026-08-13 17:00

引言:为什么要把热榜落库前几篇我们已经解决了「怎么抓」——用 httpx 监听百度金融sapi/v1/ranks接口,拿到了 A股、美股、基金三类实时热榜。但只把数据打印到控制台或存成 CSV,价值其实非常有限:无法回溯历史:热榜的精髓在于「变化」。今天涨得最猛的股票和三天前的可能完全不同,而 CSV 平铺之后,时间序列分析几乎无从下手。无法去重与对账:同一支标的在一次抓取里可能重复出现,多个榜单之间还有交集,需要一张维度表来归一。查询低效:想查「某基金过去 7 天进入热榜 Top10 的次数」,CSV 要全量扫描,而 PostgreSQL 一个索引就能秒回。所以本篇的核心目标只有一句话:设计一套能长期、可追溯、易查询的榜单数据模型,把每次抓取变成一个「快照」,把榜单里的每一项变成快照下的「记录」顺带一提:落库的前提是「抓得到、抓得稳」。下面第 3 节我会先讲如何接入亿牛云代理 IP,把采集的稳定性问题前置解决,否则后面再漂亮的表结构也只是无米之炊。一、技术选型:为什么是 SQLAlchemy + PostgreSQL

维度

选择

理由

ORM

SQLAlchemy 2.0

类型注解友好(Mapped/mapped_column),生态成熟,迁移工具 Alembic 开箱即用

数据库

PostgreSQL 15+

原生 JSONB 兜底半结构化字段、丰富索引(BRIN 适合时序)、CTE 写分析查询很顺手

驱动

psycopg(psycopg3)

SQLAlchemy 2.0 官方推荐,异步支持好

重试

tenacity

为代理请求叠加指数退避重试,提升采集健壮性

为什么不用 MySQL?榜单附带大量行情快照字段(价格、涨跌幅、成交额),后续还会追加不固定的扩展字段,PostgreSQL 的JSONB能优雅兜底,避免频繁ALTER TABLE依赖(放到pyproject.toml):二、理解ranks接口的数据结构动手建表前,先看清接口到底返回什么。百度金融sapi/v1/ranks的核心字段(已简化)大致是:关键洞察:一次请求 = 一个市场的一个榜单 = 一个快照。A股、美股、基金分别请求,所以「市场(market)」和「榜单类型(board)」是天然的维度。另外注意,价格、涨跌幅都是字符串(百度的老坑),落库前必须清洗成数值。三、采集稳定性:接入代理 IP热榜要长期采集,最大的敌人是「IP 被限频/封禁」。一旦出口 IP 被百度金融拉黑,后面整张表都无从填充。这里我用亿牛云隧道代理来兜底——它最大的好处是「每次请求自动切换出口 IP」,天然规避单 IP 的频率限制,对爬虫极其友好。亿牛云的隧道代理是一个固定的转发入口,认证信息(用户名/密码)由订单分配,建议从环境变量读取,切勿硬编码进代码库:几个工程要点:隧道 vs 动态转发:隧道代理适合「每次请求换 IP」的高频场景;若你买的是短效独享 IP(IP:PORT 形式),则需自己维护一个 IP 池并定时刷新,复杂度更高。新手直接用隧道代理最省心。代理要覆盖 http 和 https:百度金融接口是 HTTPS,但proxies字典里两个协议都配上才稳妥,避免某些 httpx 版本走直连。超时 + 重试缺一不可:代理链路比直连多一跳,10s 超时配合 3 次指数退避,能过滤掉绝大多数瞬时失败。四、数据模型设计(核心)我们采用**「维度表 + 快照表 + 记录表」**三层结构,这是时序榜单数据最稳的范式:4.1 三张表职责security(symbol, name, market, security_type)。一支标的只存一行,靠symbol去重,解决跨榜单重复问题。rank_snapshot(market, board, fetched_at, source)。代表某次抓取动作,是所有记录的根。rank_record(snapshot_id, security_id, board_rank, price, change_pct, extra)。每条热榜项一行,extra用 JSONB 兜底。4.2 SQLAlchemy 模型代码五、落库逻辑:快照 + 增量写入写入必须幂等:同一批次重复跑不能插出重复记录。策略是「先 upsert security,再建 snapshot,最后 bulk 插入 record」。

1

2

3

4

5

6

7

8

9

10

11

12

13

14

15

16

17

18

19

20

21

22

23

24

25

26

27

28

29

30

31

32

33

34

35

36

37

38

39

40

41

42

43

44

45

46

47

48

49

50

51

from sqlalchemy import select

from sqlalchemy.ext.asyncio import AsyncSession

async def save_snapshot(

session: AsyncSession,

market: str,

board: str,

items: list[dict],

) -> None:

"""把一次抓取结果落库。items 为清洗后的榜单项列表。

   Args:

       session: 异步数据库会话

       market: 市场代码 ('A'/'US'/'FUND')

       board: 榜单类型

       items: [{'symbol','name','price','change_pct','board_rank',**extra}]

   """

# 1. upsert 标的维度,保证 security 唯一

for it in items:

stmt = select(Security).where(Security.symbol == it["symbol"])

sec = (await session.execute(stmt)).scalar_one_or_none()

if sec is None:

sec = Security(

symbol=it["symbol"],

name=it["name"],

market=market,

security_type=board,

)

session.add(sec)

await session.flush()  # 拿到 security.id

# 2. 新建快照

snap = RankSnapshot(market=market, board=board)

session.add(snap)

await session.flush()

# 3. 批量插入记录

records = [

RankRecord(

snapshot_id=snap.id,

security_id=sec.id,

board_rank=it["board_rank"],

price=it["price"],

change_pct=it["change_pct"],

extra=it.get("extra", {}),

)

for it, sec in zip(items, _resolve_secs(session, items))

]

session.add_all(records)

await session.commit()

工程提示:_resolve_secs用缓存(如dict映射 symbol→id)替换上面的逐条查询,能把 N+1 查询压成 1 次IN查询,抓取量大时必做。六、常用查询示例查某基金近 7 天进 Top10 的次数:

1

2

3

4

5

6

7

8

9

10

11

12

13

14

15

16

from sqlalchemy import func as f, select

from datetime import datetime, timedelta, timezone

since = datetime.now(timezone.utc) - timedelta(days=7)

stmt = (

select(Security.symbol, Security.name, f.count().label("top10_cnt"))

.join(RankRecord, RankRecord.security_id == Security.id)

.join(RankSnapshot, RankSnapshot.id == RankRecord.snapshot_id)

.where(

RankSnapshot.market == "FUND",

RankRecord.board_rank <= 10,

RankSnapshot.fetched_at >= since,

)

.group_by(Security.symbol, Security.name)

.order_by(f.count().desc())

)

查某标的历次榜单名次的走势(时序):

1

2

3

4

5

6

7

stmt = (

select(RankSnapshot.fetched_at, RankRecord.board_rank, RankRecord.price)

.join(RankRecord, RankRecord.snapshot_id == RankSnapshot.id)

.join(Security, Security.id == RankRecord.security_id)

.where(Security.symbol == "SH600519")

.order_by(RankSnapshot.fetched_at)

)

七、踩坑与优化建议价格/涨跌幅是字符串:百度返回"+2.31%",必须正则提取数字并转float,否则排序全错。代理是采集的命门:直连极易被限频,生产务必走亿牛云这类隧道代理,并把认证信息放环境变量,密钥绝不入库。快照时间用服务端now():爬虫机器时钟不可信,统一用 PG 的server_default,分析时才对齐。BRIN 索引救大表rank_record会极速膨胀,对snapshot_id建 BRIN 索引,比 B-tree 省空间且时序范围查询更快。软删除 + 分区:按月对rank_record做声明式分区,冷数据可快速 detach 归档。别把 extra 当主字段用:JSONB 只兜底,高频查询字段一定要提升为独立列并加索引。八、小结本篇给出了百度金融热榜的完整落库链路:先用亿牛云隧道代理保证采集稳定,再以 security 维度表做去重、rank_snapshot 记抓取动作、rank_record 存明细,配合 SQLAlchemy 2.0 的类型安全模型与 PostgreSQL 的 JSONB/BRIN 能力。这套结构既能支撑简单的历史回溯,也能平滑演进到跨市场热度分析。下一篇我们会用schedule+Typer把它变成每天自动跑的定时任务。运行环境:Python 3.11+ / SQLAlchemy 2.0 / PostgreSQL 15+。建表用Base.metadata.create_all(engine)或 Alembic 迁移,生产环境强烈建议后者。亿牛云代理需在.env中配置YINIU_USERYINIU_PASS


评论