ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

从零搭建自托管金融数据服务:架构、抓取与存储实战

从零搭建自托管金融数据服务:架构、抓取与存储实战 1. 金融数据服务从零搭建的核心思路拆解1.1 为什么我要自己动手做一套金融数据服务先说清楚这套东西是干什么的。financial-services这个名字听起来很泛实际上我做的是一套面向个人开发者和小型团队的自托管金融数据聚合与分发服务。它能做什么简单讲就是把行情数据、财报数据、宏观经济指标这几类信息从多个公开数据源抓取、清洗、标准化之后通过统一的接口对外提供查询能力。解决的问题也很直接市面上成熟的金融数据API要么贵得离谱要么免费额度少得可怜要么字段定义混乱到让人抓狂。适合谁来参考有一定后端基础、想自己掌控数据管道的开发者或者做量化研究、个人记账工具、投资看板这类小项目的朋友。我最初动这个念头是因为手上有个小工具需要每天拉取几十只标的的收盘数据做统计。用了几家免费接口之后发现要么限流卡得死死的要么字段隔三差五变要么历史数据缺胳膊少腿。折腾了两个月我决定干脆自己搭一套把数据抓下来存到自己库里想怎么查就怎么查。这套服务跑了大半年日均处理几万条记录稳定性还不错所以把设计思路和踩过的坑整理出来。核心关键词就一个financial-services。但围绕它展开的东西不少——数据源管理、抓取调度、字段映射、存储选型、接口设计、容错重试每一块都有讲究。下面我按实际搭建的顺序把每个环节的选择理由和操作细节讲透。1.2 整体架构选型为什么是抓取-清洗-存储-服务四层我见过不少人一上来就想搞微服务、上消息队列、拆十几个模块结果光环境搭建就耗掉一周最后业务逻辑没写几行。我的建议是个人项目先把链路跑通再谈优化。所以这套financial-services我采用的是最朴素的四层结构抓取层负责从各个数据源拉取原始数据处理请求频率、重试、超时。清洗层把不同来源的字段统一成一套标准模型处理缺失值、异常值、单位换算。存储层按数据类型分别落到关系库和时序库兼顾查询灵活性和写入性能。服务层对外暴露 REST 接口做参数校验、缓存、限流。为什么这么分因为金融数据的痛点在于来源异构。同一个收盘价A 源叫closeB 源叫closing_priceC 源可能给你一个字符串还带货币符号。如果不把清洗单独抽一层后面每加一个数据源就要改一遍业务代码维护成本爆炸。分层之后新增数据源只需要在抓取层加一个适配器在清洗层加一套映射规则服务层完全不用动。技术栈方面我选的是 Python 做抓取和清洗生态成熟处理数据方便PostgreSQL 存结构化数据财报、标的元信息TimescaleDB 存时序行情它是 PostgreSQL 的扩展不用额外维护一套数据库FastAPI 做服务层自带文档性能足够。这套组合的好处是只用维护一个数据库实例备份、迁移都省心。如果你量级更大可以把时序部分换成专门的列式存储但个人项目真没必要。提示不要一上来就追求高可用分布式。个人项目的瓶颈通常在数据源限流而不是你的服务性能。先把单机跑稳比什么都强。1.3 数据源的选择逻辑与合规边界选数据源这件事我的原则是优先官方、其次社区、最后才考虑爬取。官方接口通常有明确的字段定义和使用条款稳定性最好社区维护的开源数据包更新及时但要注意版本兼容爬取是最后手段因为页面结构一变你就得跟着改而且频率控制不好容易被封。具体到financial-services我用到的数据源分三类数据类型来源类型更新频率注意事项日线行情公开接口每日收盘后注意复权处理前复权后复权要分清财务报表公开披露季度字段口径可能调整需保留原始快照宏观指标统计发布月度/季度发布时间不固定需轮询检测这里要特别强调合规边界。所有数据我只用于个人研究和小范围非商业用途抓取频率严格控制在对方允许的范围内并且遵守 robots 协议。如果你打算商用务必去谈正式授权不要心存侥幸。这不是技术问题是底线问题。另外数据源一定要做冗余。我早期只依赖一个源结果对方维护了三天我的服务直接断粮。后来改成主备双源主源失败自动切备源虽然字段要额外做一层对齐但稳定性提升非常明显。2. 核心细节解析与实操要点2.1 字段标准化一套模型打天下的关键金融数据最烦人的地方就是字段不统一。我举个真实例子同样是成交量有的源给你股数有的源给你手数1手100股还有的源在指数数据里给你的是金额而不是股数。如果不做标准化你算出来的指标全是错的。我的做法是定义一套内部标准模型所有外部数据进来先转成这个模型。以日线行情为例核心字段定义如下# 标准行情模型简化版 class StandardBar: symbol: str # 统一格式的标的代码如 600519.SH trade_date: date # 交易日期 open: Decimal # 开盘价统一为标的计价货币 high: Decimal # 最高价 low: Decimal # 最低价 close: Decimal # 收盘价 volume: int # 成交量统一为股 amount: Decimal # 成交额统一为计价货币 adjust_flag: str # 复权标志none/qfq/hfq关键点在于单位统一和代码统一。代码统一我用的规则是数字代码.交易所后缀交易所后缀用 SH、SZ、HK、US 这类通用标识。单位统一则要在清洗层做换算比如手数乘以 100 转成股数。这些换算规则我全部写在配置文件里而不是硬编码在代码中方便后续调整。注意复权处理一定要在入库前想清楚。我的做法是同时存原始价格和复权因子查询时按需计算。这样既保留了原始数据又能灵活支持前复权、后复权两种需求。如果只存一种复权价后面想换算法就得重新拉全部历史数据非常痛苦。2.2 抓取调度别让限流把你打趴下抓取层最容易出问题的地方就是频率控制。我踩过最大的坑是一开始用多线程并发抓取觉得速度快结果十分钟就被对方限流IP 被封了整整一天。后来老老实实改成令牌桶限流把请求速率压到对方允许的阈值以下反而再没出过问题。我的调度设计是这样的每个数据源维护一个独立的令牌桶速率和突发量按对方文档配置。抓取任务用队列串行化避免并发打满。每个请求设置超时和重试重试用指数退避最多三次。失败任务进入死信队列人工介入或定时重跑。令牌桶的实现我用的是简单的计数加时间窗口没必要上 Redis 那么重。核心逻辑就是每次请求前检查当前窗口内已用令牌数超了就 sleep 到下一个窗口。代码大概长这样import time from threading import Lock class TokenBucket: def __init__(self, rate: float, capacity: int): self.rate rate # 每秒补充的令牌数 self.capacity capacity # 桶容量 self.tokens capacity self.last_time time.time() self.lock Lock() def acquire(self, n: int 1): with self.lock: now time.time() # 按时间差补充令牌 self.tokens min( self.capacity, self.tokens (now - self.last_time) * self.rate ) self.last_time now if self.tokens n: self.tokens - n return True return False调用的时候如果acquire返回 False就 sleep 一小段时间再试。这个方案简单可靠实测下来很稳。2.3 存储选型关系库和时序库怎么分工存储这块我纠结了很久最后定下来的方案是分工存储PostgreSQL 普通表存标的元信息、财报数据、宏观指标。这些数据更新频率低但查询维度多关系库的灵活查询能力正好合适。TimescaleDB 超表存日线、分钟线这类时序数据。它自动按时间分区查询某段时间的数据非常快而且支持连续聚合算均线、算波动率都很方便。为什么不用一个库全搞定因为时序数据的写入模式是高频追加而财报数据的写入模式是低频更新两者的索引策略和存储优化方向完全不同。混在一起会导致两边都做不好。TimescaleDB 的好处是它本身就是 PostgreSQL 扩展不用额外维护一套数据库运维成本几乎为零。建表的时候有个细节要注意时序表的主键设计。我用的是(symbol, trade_date)联合主键配合时间分区。这样既能保证同一标的同一天不重复又能让按时间范围的查询走分区裁剪速度很快。如果只按时间分区不设标的索引查单只标的的历史数据会全表扫描慢得让人怀疑人生。3. 实操过程与核心环节实现3.1 从零搭建环境准备与依赖安装先说环境。我用的是 Ubuntu 22.04Python 3.11。数据库部分PostgreSQL 15 加上 TimescaleDB 2.x 扩展。整个搭建过程我整理成了可复现的步骤你照着做基本不会出问题。第一步装数据库。用官方源装 PostgreSQL然后加 TimescaleDB 的源装扩展# 安装 PostgreSQL sudo apt install postgresql-15 postgresql-client-15 # 添加 TimescaleDB 源并安装 sudo add-apt-repository ppa:timescale/timescaledb-ppa sudo apt update sudo apt install timescaledb-2-postgresql-15 # 配置并重启 sudo timescaledb-tune --quiet --yes sudo systemctl restart postgresql第二步建库建表。我习惯把建表语句写成迁移脚本方便版本管理。核心的几张表-- 标的元信息表 CREATE TABLE symbols ( symbol VARCHAR(20) PRIMARY KEY, name VARCHAR(100) NOT NULL, exchange VARCHAR(10) NOT NULL, asset_type VARCHAR(20) NOT NULL, listed_date DATE, created_at TIMESTAMPTZ DEFAULT NOW() ); -- 日线行情表TimescaleDB 超表 CREATE TABLE daily_bars ( symbol VARCHAR(20) NOT NULL, trade_date DATE NOT NULL, open NUMERIC(18,4), high NUMERIC(18,4), low NUMERIC(18,4), close NUMERIC(18,4), volume BIGINT, amount NUMERIC(20,4), adjust_flag VARCHAR(4) DEFAULT none, PRIMARY KEY (symbol, trade_date, adjust_flag) ); -- 转成超表按交易日期分区 SELECT create_hypertable(daily_bars, trade_date, chunk_time_interval INTERVAL 1 year);第三步装 Python 依赖。抓取用httpx比 requests 更适合异步数据处理用pandas数据库操作用sqlalchemy加psycopg2服务层用fastapi加uvicornpip install httpx pandas sqlalchemy psycopg2-binary fastapi uvicorn pydantic提示psycopg2-binary在部分系统上装不上可以换成psycopg2源码编译或者用asyncpg走异步。我实测psycopg2-binary在 Ubuntu 上没问题但如果你用 Alpine 镜像记得装postgresql-dev。3.2 抓取适配器的编写与字段映射抓取层的核心是适配器模式。每个数据源写一个适配器类实现统一的接口fetch负责拉数据parse负责解析成标准模型。这样新增数据源只需要加一个类不用动其他代码。我以某公开行情接口为例写一个简化的适配器from abc import ABC, abstractmethod from decimal import Decimal from datetime import date class BaseAdapter(ABC): abstractmethod def fetch(self, symbol: str, start: date, end: date) - list[dict]: 拉取原始数据 pass abstractmethod def parse(self, raw: dict) - dict: 解析成标准模型 pass class SourceAAdapter(BaseAdapter): def __init__(self, client, bucket): self.client client self.bucket bucket def fetch(self, symbol, start, end): while not self.bucket.acquire(): time.sleep(0.1) resp self.client.get( /api/daily, params{code: symbol, from: start.isoformat(), to: end.isoformat()}, timeout10 ) resp.raise_for_status() return resp.json()[data] def parse(self, raw): # 关键单位换算和字段映射 return { symbol: self._normalize_symbol(raw[code]), trade_date: date.fromisoformat(raw[date]), open: Decimal(str(raw[open])), high: Decimal(str(raw[high])), low: Decimal(str(raw[low])), close: Decimal(str(raw[close])), volume: int(raw[vol]) * 100, # 手转股 amount: Decimal(str(raw[turnover])), adjust_flag: none }这里有个细节值得展开为什么用 Decimal 而不是 float。金融数据对精度极其敏感float 的浮点误差在累加计算时会放大。比如你算一年的累计收益用 float 可能差出好几个基点。Decimal 虽然慢一点但精度可靠这是金融场景的硬要求。字段映射我建议单独抽成配置文件比如 YAMLsource_a: symbol: code trade_date: date open: open high: high low: low close: close volume: vol amount: turnover volume_multiplier: 100这样换数据源或者对方改字段名改配置就行不用动代码。3.3 清洗入库的完整流程与幂等处理清洗入库这一步核心要求是幂等。什么意思就是同一批数据重复跑结果不能变。金融数据经常需要补拉、重跑如果幂等做不好数据就乱了。我的做法是用 PostgreSQL 的ON CONFLICT语法做 upsertfrom sqlalchemy.dialects.postgresql import insert def upsert_bars(session, bars: list[dict]): if not bars: return 0 stmt insert(DailyBar).values(bars) stmt stmt.on_conflict_do_update( index_elements[symbol, trade_date, adjust_flag], set_{ open: stmt.excluded.open, high: stmt.excluded.high, low: stmt.excluded.low, close: stmt.excluded.close, volume: stmt.excluded.volume, amount: stmt.excluded.amount, } ) result session.execute(stmt) session.commit() return result.rowcount这样无论跑多少次最终数据都是一致的。配合批量提交每 1000 条提交一次写入性能也不错。我实测单机每秒能写两三万条对个人项目完全够用。清洗环节还有几个必做的检查空值检查关键字段开高低收不能为空为空直接丢弃并记录日志。逻辑校验最高价必须大于等于最低价开盘价和收盘价要落在高低区间内否则标记为异常。重复检查同一标的同一日期同一复权标志只能有一条靠主键约束保证。时间连续性检查交易日之间不应该有非交易日的空洞发现异常要排查数据源。注意异常数据不要直接删要存到单独的异常表里。我早期直接丢弃后来发现某些异常其实是数据源的特殊处理比如停牌日的价格填充删了就找不回来了。存异常表人工复核才是稳妥做法。4. 常见问题与排查技巧实录4.1 数据源限流与封禁的应对策略这是最高频的问题。表现是请求突然开始返回 429 或者直接超时严重的话 IP 被临时封禁。排查思路分三步第一确认是不是限流。看返回状态码和响应头很多接口会在 header 里告诉你剩余额度。如果状态码是 429基本可以确定。第二检查自己的请求速率。我写了个简单的统计脚本记录每分钟的实际请求数跟配置的令牌桶速率对比。经常发现是某个定时任务和主任务撞车了导致瞬时速率超标。第三降速并加退避。把令牌桶速率调低重试间隔用指数退避。如果已经被封就等封禁期过了再恢复期间切到备用源。我整理了一个速查表现象可能原因处理方式429 状态码请求速率超限降低令牌桶速率加退避重试连接超时网络抖动或对方限流增加超时时间切换备用源返回空数据参数错误或数据未更新检查参数确认数据发布时间字段缺失对方接口变更对比文档更新字段映射配置数据明显错误单位或口径变化核对原始数据更新换算规则4.2 数据质量问题的排查方法数据质量问题最隐蔽因为服务不报错但算出来的结果是错的。我遇到过几次典型情况一次是某只标的的成交量突然比前一天大了 100 倍排查发现是数据源把单位从手改成了股而我的换算规则还在乘以 100。解决办法是在清洗层加量级校验如果某天的成交量偏离近 20 日均值超过 10 倍就标记为可疑人工复核。另一次是财报数据的字段口径变了同一个营业收入字段新版本包含了子公司数据旧版本没有。这种问题很难自动发现我的做法是保留原始快照每次抓取都把原始 JSON 存一份到对象存储出问题时可以回溯对比。排查数据质量问题我的经验是建立一套校验规则库每次入库前跑一遍价格区间校验价格不能为负不能超过合理上限。量价关系校验成交额除以成交量应该约等于均价偏差过大要查。时间序列校验相邻交易日的数据变化率不能超过阈值。跨源校验主备源的数据应该一致不一致要告警。4.3 服务层的性能优化与缓存设计服务层上线初期我发现某些查询特别慢比如查某只标的近五年的日线。排查下来是两个原因一是没走分区裁剪二是每次都实时查库。优化方案分两步。第一步确保查询走索引和分区。TimescaleDB 的超表在按时间范围查询时会自动裁剪分区但前提是你的查询条件里要有时间字段。我强制要求所有行情查询必须带时间范围不允许全表扫描。第二步加缓存。金融数据的特性是读多写少历史数据几乎不变非常适合缓存。我用的是两级缓存内存缓存进程内 LRU加 Redis跨进程共享。缓存键设计成bars:{symbol}:{start}:{end}:{adjust}过期时间按数据新鲜度设置——历史数据缓存一天当天数据缓存五分钟。from functools import lru_cache import redis import json r redis.Redis(hostlocalhost, port6379, db0) def get_bars(symbol, start, end, adjustnone): cache_key fbars:{symbol}:{start}:{end}:{adjust} # 先查 Redis cached r.get(cache_key) if cached: return json.loads(cached) # 查库 bars query_from_db(symbol, start, end, adjust) # 写缓存历史数据缓存久一点 ttl 300 if end date.today().isoformat() else 86400 r.setex(cache_key, ttl, json.dumps(bars, defaultstr)) return bars实测下来加了缓存之后热门查询的响应时间从几百毫秒降到几毫秒数据库压力也小了很多。4.4 定时任务的可靠性保障financial-services里有很多定时任务每日收盘后拉行情、季度拉财报、月度拉宏观。这些任务最怕的是静默失败——任务跑了但没拉到数据或者根本没跑而你浑然不知。我的保障措施有三层第一层任务状态记录。每次任务执行都往数据库写一条记录包含开始时间、结束时间、处理条数、状态。这样一眼就能看出哪个任务没跑或者跑失败了。第二层数据新鲜度监控。写一个检查脚本每天定时检查各类数据的最新日期。如果行情数据的最新日期不是最近一个交易日就发告警。这个比任务状态更可靠因为它直接检查结果。第三层失败自动重试加人工兜底。任务失败自动重试三次三次都失败就发通知我手动介入。通知我用的是邮件加一个简单的 webhook不依赖第三方服务避免额外故障点。提示定时任务千万别用cron直接调 Python 脚本因为环境变量和路径经常对不上。我推荐用systemd timer或者APScheduler前者适合系统级任务后者适合应用内调度。我用的是 APScheduler配置写在代码里版本管理方便。5. 接口设计与对外服务能力5.1 REST 接口的字段设计与版本管理服务层的接口设计我的原则是稳定优先。金融数据的使用者包括我自己最怕接口字段变来变去。所以我在 URL 里加了版本号比如/api/v1/bars任何破坏性变更都开新版本老版本至少保留半年。核心接口就几个接口方法说明关键参数/api/v1/symbolsGET查询标的列表exchange, asset_type/api/v1/barsGET查询行情symbol, start, end, adjust/api/v1/financialsGET查询财报symbol, period/api/v1/macroGET查询宏观指标indicator, start, end返回格式统一用 JSON字段名用下划线风格。分页用limit和offset默认 limit 是 100最大 1000。时间字段统一用 ISO 8601 格式避免时区歧义。from fastapi import FastAPI, Query from datetime import date app FastAPI(titleFinancial Services API, version1.0) app.get(/api/v1/bars) def get_bars( symbol: str Query(..., description标的代码), start: date Query(..., description开始日期), end: date Query(..., description结束日期), adjust: str Query(none, regex^(none|qfq|hfq)$), limit: int Query(100, le1000) ): bars service.query_bars(symbol, start, end, adjust, limit) return {code: 0, data: bars, count: len(bars)}参数校验交给 FastAPI 的 Pydantic 模型省心又可靠。错误返回统一格式code非零表示出错message给出原因。5.2 限流、鉴权与使用统计虽然是个人项目但既然对外提供服务基本的防护还是要做。限流我用的是基于 IP 的滑动窗口每个 IP 每分钟最多 60 次请求。鉴权用简单的 API Key每个使用者分配一个 key记录在数据库里。使用统计这块我记录每个 key 的调用次数、调用时间、查询参数。一方面是防止滥用另一方面也能看出哪些数据最受欢迎指导后续优化方向。统计信息每天汇总一次存到单独的统计表里。# 简单的滑动窗口限流 from collections import defaultdict, deque import time class RateLimiter: def __init__(self, max_requests60, window60): self.max_requests max_requests self.window window self.records defaultdict(deque) def allow(self, key: str) - bool: now time.time() dq self.records[key] # 移除窗口外的记录 while dq and dq[0] now - self.window: dq.popleft() if len(dq) self.max_requests: return False dq.append(now) return True这个实现简单单机够用。如果要多实例部署把记录换成 Redis 就行。5.3 数据导出与批量查询的优化有些场景需要批量拉取大量数据比如回测要拉几百只标的的历史行情。如果一只一只查接口调用次数太多效率低。我专门做了一个批量接口/api/v1/bars/batch支持一次传多个标的代码内部并行查询后合并返回。批量查询的优化点在于减少数据库往返。我用的是WHERE symbol IN (...)加时间范围一次查出来再按标的分组。配合缓存批量查询的响应时间能控制在秒级。导出功能支持 CSV 和 JSON 两种格式通过format参数指定。CSV 适合导入 Excel 做分析JSON 适合程序处理。导出大文件时用流式响应避免内存爆掉。from fastapi.responses import StreamingResponse import csv import io app.get(/api/v1/bars/export) def export_bars(symbol: str, start: date, end: date, format: str csv): bars service.query_bars(symbol, start, end, none, limit100000) if format csv: output io.StringIO() writer csv.DictWriter(output, fieldnamesbars[0].keys()) writer.writeheader() writer.writerows(bars) output.seek(0) return StreamingResponse( iter([output.getvalue()]), media_typetext/csv, headers{Content-Disposition: fattachment; filename{symbol}.csv} ) return {code: 0, data: bars}6. 运维监控与长期演进6.1 日志、指标与告警的最小可用组合个人项目的运维不需要搞得太复杂但日志、指标、告警这三样必须有。我的方案是日志用 Python 标准库的 logging输出到文件并按天轮转。关键操作抓取、入库、接口调用都记日志格式统一方便 grep。指标用 Prometheus 的 Python 客户端暴露几个核心指标——抓取成功数、失败数、入库条数、接口响应时间。Prometheus 加 Grafana 的部署成本不高但可视化效果很好。告警基于指标设置阈值比如连续 3 次抓取失败或数据新鲜度超过 24 小时就触发告警。告警渠道用邮件简单可靠。日志格式我统一成 JSON方便后续用工具解析import logging import json class JsonFormatter(logging.Formatter): def format(self, record): log_obj { time: self.formatTime(record), level: record.levelname, module: record.module, message: record.getMessage() } if record.exc_info: log_obj[exception] self.formatException(record.exc_info) return json.dumps(log_obj, ensure_asciiFalse) logger logging.getLogger(financial-services) handler logging.FileHandler(app.log) handler.setFormatter(JsonFormatter()) logger.addHandler(handler) logger.setLevel(logging.INFO)6.2 数据备份与恢复的实操方案金融数据是长期积累的资产丢了很难补回来。备份策略我采用每日全量加实时增量每日凌晨用pg_dump做全量备份压缩后存到本地和对象存储各一份。关键表开启 WAL 归档支持时间点恢复。备份文件保留最近 30 天更早的按月归档。恢复演练我每季度做一次确保备份真的能用。我踩过的坑是备份脚本跑了半年结果恢复时发现某个表没被包含进去。所以备份脚本一定要显式列出所有表不要用通配符并且定期验证。#!/bin/bash # 每日备份脚本 BACKUP_DIR/data/backup DATE$(date %Y%m%d) pg_dump -Fc -d financial_services -f ${BACKUP_DIR}/fs_${DATE}.dump # 上传到对象存储用你习惯的工具 # 清理 30 天前的备份 find ${BACKUP_DIR} -name fs_*.dump -mtime 30 -delete6.3 后续可扩展的方向与个人建议这套financial-services跑到现在我觉得还有几个方向可以继续打磨。一是增加技术指标计算把均线、MACD、RSI 这些常用指标在服务层算好缓存起来省得每个使用者自己算。二是支持实时行情推送用 WebSocket 把盘中数据推给订阅者不过这需要更严格的数据源授权。三是做数据质量看板把前面提到的校验规则结果可视化一眼看出哪些数据有问题。不过我的建议是别贪多。个人项目最容易死在什么都想做上。先把核心的行情和财报数据做扎实接口稳定数据准确这已经能覆盖 80% 的需求了。剩下的功能等真正有需求再加。最后分享一个我踩过的坑早期我为了追求实时把抓取频率设得很高结果不仅被限流还因为数据源本身有延迟抓到的实时数据其实是过时的反而误导了判断。后来想明白了金融数据服务的关键不是快而是准和稳。宁可晚几分钟也要保证数据是对的。这个认知转变比任何技术优化都重要。
返回列表