1. 金融数据服务从零搭建的完整思路
1.1 为什么我要自己搭一套金融数据服务
先说清楚这个项目到底在干什么。financial-services这个名字听起来很泛,实际上我做的事情是:搭建一套能稳定拉取、清洗、存储、对外输出金融行情与基础面数据的后端服务。它解决的核心问题是——市面上现成的数据接口要么贵得离谱,要么限流严重,要么字段残缺还得自己拼,而量化研究、行情看板、策略回测这些场景又对数据的时效性、完整性、一致性有硬性要求。
这套东西适合谁?如果你正在做量化策略验证、搭建个人行情监控面板、给内部团队做数据中台,或者单纯想搞清楚金融数据从交易所到你自己数据库这条链路是怎么跑通的,那这篇内容就是写给你的。不需要你是金融科班出身,但至少要能看懂 JSON、会写基本的 SQL、对 HTTP 接口调用不陌生。
我踩过的第一个坑就是:一开始觉得“不就是调个接口存数据库吗”,结果真上手才发现,数据源的时间戳格式五花八门、复权因子对不上、停牌日的空值处理逻辑各家不一样,光是统一字段就花了我整整一周。所以这套服务的价值不在于“能拉到数据”,而在于把脏活累活封装掉,对外只暴露干净的、语义统一的数据。
1.2 整体架构怎么设计才不返工
我的架构分四层,从下往上依次是:数据采集层、数据清洗层、存储层、服务接口层。这个分层不是拍脑袋定的,而是根据“数据流向”和“故障隔离”两个原则来的。
采集层负责跟外部数据源打交道,不管是公开的行情接口、CSV 文件导入,还是第三方付费源,全部在这一层做适配。为什么要单独抽一层?因为数据源是会变的——今天用的源明天可能改字段、后天可能涨价,如果采集逻辑和业务逻辑混在一起,换个源就得改半个项目。采集层对外只输出一种内部标准格式,我叫它RawTick,包含symbol、timestamp、open/high/low/close、volume、source这几个必填字段。
清洗层是整套服务最核心也最容易被低估的部分。它干三件事:时间对齐、缺失值处理、异常值过滤。时间对齐是指把不同源的数据统一到同一时区、同一粒度(比如都归到分钟线);缺失值处理要区分“停牌导致的真缺失”和“采集失败导致的假缺失”,前者填 NULL 并标记,后者要触发重试;异常值过滤则是把明显错误的价格(比如某只股票突然出现 0.01 元的价格)拦下来。
存储层我选了PostgreSQL + TimescaleDB 扩展。为什么不用纯时序数据库比如 InfluxDB?因为金融数据不只是时序,还有大量的关系型查询需求——比如“找出所有市盈率低于 15 且近三日成交量放大的股票”,这种多条件联合查询在关系型库里写起来舒服得多。TimescaleDB 的好处是保留了 PostgreSQL 的全部能力,同时给时序场景做了分区和压缩优化,实测下来单表存几亿条 K 线数据,查询延迟依然能控制在毫秒级。
服务接口层用FastAPI暴露 REST 接口,同时留了 WebSocket 通道推实时行情。选 FastAPI 而不是 Flask 或 Django,主要看中它的异步支持和自动生成的 OpenAPI 文档——后者在团队协作时省了大量写接口文档的时间。
提示:架构分层的关键原则是“上层不知道下层的实现细节”。采集层换源,清洗层和存储层完全无感知;存储层从 PostgreSQL 换成 ClickHouse,接口层的代码一行不用改。这个约束在项目初期就要定死,否则后期重构成本极高。
1.3 技术选型背后的取舍逻辑
有人可能会问:为什么不用现成的量化框架比如 Backtrader 或 Zipline 自带的数据模块?我的理由是——那些框架的数据模块是为回测服务的,字段和粒度都是固定的,而我的需求是对外提供通用数据服务,要支持多种粒度、多种资产类别、多种查询模式,框架自带的数据层根本不够用。
数据库选型上我对比过三个方案:
| 方案 | 优势 | 劣势 | 适用场景 |
|---|---|---|---|
| PostgreSQL + TimescaleDB | 关系型查询强、生态成熟、支持 JSON 字段 | 写入吞吐不如专用时序库 | 中低频数据、多条件查询 |
| InfluxDB | 写入极快、压缩率高 | 关系型查询弱、SQL 支持有限 | 纯时序监控场景 |
| ClickHouse | 列式存储、聚合查询极快 | 运维复杂度高、事务支持弱 | 海量数据分析 |
最终选 PostgreSQL + TimescaleDB,是因为我的查询模式里“多条件筛选 + 排序 + 分页”占比超过 60%,这类查询在关系型库里写起来最顺手。写入方面,通过批量插入和连接池优化,实测单机每秒能稳定写入 5 万条 K 线记录,对个人和小团队场景完全够用。
2. 数据采集与清洗的核心细节
2.1 采集频率与重试机制怎么定
采集频率不是拍脑袋定的,要根据数据源的更新节奏和你的使用场景来算。日线数据每天收盘后拉一次就够,分钟线数据如果做日内策略,至少得每分钟拉一次。我目前的配置是:日线数据每天 15:30 触发全量拉取,分钟线数据每 60 秒增量拉取一次。
重试机制我用了指数退避 + 最大重试次数的组合。具体参数是:首次失败后等 2 秒重试,第二次等 4 秒,第三次等 8 秒,最多重试 5 次。为什么用指数退避而不是固定间隔?因为数据源限流通常是“短时间内请求过多”导致的,固定间隔重试容易撞上限流窗口,指数退避能给数据源足够的恢复时间。
import time import requests def fetch_with_retry(url, params, max_retries=5): for attempt in range(max_retries): try: resp = requests.get(url, params=params, timeout=10) if resp.status_code == 200: return resp.json() elif resp.status_code == 429: wait = 2 ** attempt time.sleep(wait) else: resp.raise_for_status() except requests.RequestException as e: if attempt == max_retries - 1: raise time.sleep(2 ** attempt) return None这段代码里有个细节:429状态码单独处理,因为它是明确的限流信号,其他错误码走通用重试逻辑。另外timeout=10是必须的,我遇到过数据源不响应但也不断开连接的情况,没有超时设置的话线程会一直挂在那里。
2.2 字段标准化与复权处理
不同数据源返回的字段名和格式差异极大。有的用trade_date,有的用date;有的价格是字符串,有的是浮点数;有的成交量单位是股,有的是手。我的做法是在采集层和清洗层之间加一个字段映射表,用配置的方式做转换,而不是硬编码在代码里。
FIELD_MAPPING = { "source_a": { "trade_date": ("timestamp", lambda x: f"{x[:4]}-{x[4:6]}-{x[6:]}"), "vol": ("volume", lambda x: float(x) * 100), "close": ("close", float), }, "source_b": { "date": ("timestamp", str), "volume": ("volume", float), "close_price": ("close", float), } }复权处理是金融数据里绕不开的坎。前复权、后复权、不复权三种模式,用错了会导致回测结果完全失真。我的处理逻辑是:原始数据只存不复权价格,复权因子单独存一张表,查询时根据参数动态计算。这样做的好处是数据只存一份,复权方式可以随时切换,不用重新拉数据。
复权因子的计算基于除权除息日的价格调整比例。假设某股票除权前收盘价 10 元,除权后开盘参考价 9 元,那么复权因子就是 10/9 ≈ 1.111。后复权价格 = 不复权价格 × 累计复权因子,前复权价格 = 不复权价格 × (当前复权因子 / 历史复权因子)。这个计算过程我封装成了一个 PostgreSQL 函数,查询时直接调用。
2.3 缺失值与异常值的处理策略
缺失值分两种:真缺失和假缺失。真缺失是指股票停牌、未上市、已退市等客观原因导致的数据不存在,这种填 NULL 并打上标记就行。假缺失是指采集失败、网络抖动导致的数据丢失,这种必须触发补采。
我的判断逻辑是:如果某只股票在某个交易日没有数据,先查该股票是否处于停牌状态(停牌信息从另一个接口获取),如果是停牌,标记为SUSPENDED;如果不是停牌,标记为MISSING并加入补采队列。
异常值过滤我用的是基于波动率的动态阈值。具体做法是:计算过去 20 个交易日的收益率标准差,如果当日收益率超过 5 倍标准差,就标记为疑似异常。为什么用 5 倍而不是 3 倍?因为金融数据本身波动就大,3 倍标准差会误杀很多正常的大涨大跌,5 倍是一个经验上比较平衡的值。
-- 异常值检测查询 WITH stats AS ( SELECT symbol, STDDEV(return) AS std_return FROM daily_returns WHERE timestamp >= NOW() - INTERVAL '20 days' GROUP BY symbol ) SELECT r.symbol, r.timestamp, r.return FROM daily_returns r JOIN stats s ON r.symbol = s.symbol WHERE ABS(r.return) > 5 * s.std_return;注意:异常值不要直接删除,而是标记后保留。我吃过亏——有一次某只股票真的发生了极端行情,被我的过滤逻辑删掉了,导致回测结果偏乐观。后来改成“标记但不删除”,查询时可以选择是否排除。
3. 存储层设计与查询优化实操
3.1 表结构设计与分区策略
核心表有三张:instruments(标的元信息)、daily_bars(日线数据)、minute_bars(分钟线数据)。instruments存股票代码、名称、上市日期、退市日期、所属行业等静态信息;daily_bars和minute_bars存行情数据。
分区策略上,daily_bars按年份做范围分区,minute_bars按月份做范围分区。为什么粒度不同?因为分钟线数据量是日线的几百倍,按月分区能让单个分区的数据量控制在合理范围内,查询时分区裁剪的效果更明显。
-- 创建日线表并按年分区 CREATE TABLE daily_bars ( symbol VARCHAR(10) NOT NULL, timestamp DATE NOT NULL, open NUMERIC(12,4), high NUMERIC(12,4), low NUMERIC(12,4), close NUMERIC(12,4), volume BIGINT, PRIMARY KEY (symbol, timestamp) ) PARTITION BY RANGE (timestamp); CREATE TABLE daily_bars_2024 PARTITION OF daily_bars FOR VALUES FROM ('2024-01-01') TO ('2025-01-01');索引方面,(symbol, timestamp)的联合主键已经覆盖了大部分查询场景。额外加了一个(timestamp)的单列索引,用于“查询某一天所有股票”的场景。索引不是越多越好,每加一个索引都会拖慢写入速度,我实测下来这两个索引是性价比最高的组合。
3.2 批量写入与连接池配置
写入性能是这套服务的瓶颈之一。单条 INSERT 的效率极低,我改成了批量插入 + COPY的方式。对于大批量历史数据导入,直接用 PostgreSQL 的COPY命令,速度比 INSERT 快一个数量级。
import psycopg2 from io import StringIO def bulk_insert_bars(conn, bars): buf = StringIO() for bar in bars: buf.write(f"{bar['symbol']}\t{bar['timestamp']}\t{bar['open']}\t" f"{bar['high']}\t{bar['low']}\t{bar['close']}\t{bar['volume']}\n") buf.seek(0) with conn.cursor() as cur: cur.copy_from(buf, 'daily_bars', columns=('symbol', 'timestamp', 'open', 'high', 'low', 'close', 'volume')) conn.commit()连接池我用的是psycopg2.pool.ThreadedConnectionPool,最小连接数 5,最大 20。为什么是 20?因为我的采集并发度是 10,每个采集任务需要一个连接,加上接口层的查询需求,20 个连接刚好够用且不会把数据库压垮。连接数不是越多越好,PostgreSQL 每个连接都会占用内存,超过实际需求反而会降低整体吞吐。
3.3 查询性能优化的几个关键手段
第一个手段是物化视图。对于“最新收盘价”这种高频查询但计算量不大的需求,我建了一个物化视图,每 5 分钟刷新一次。查询直接走物化视图,避免每次都去扫全表。
CREATE MATERIALIZED VIEW latest_close AS SELECT DISTINCT ON (symbol) symbol, timestamp, close FROM daily_bars ORDER BY symbol, timestamp DESC; CREATE INDEX ON latest_close (symbol);第二个手段是查询缓存。对于不常变的数据(比如历史日线),我在接口层加了 Redis 缓存,TTL 设为 1 小时。实测下来,缓存命中率超过 80%,数据库压力下降明显。
第三个手段是分区裁剪。查询时一定要带上时间范围条件,这样 PostgreSQL 才能自动裁剪掉不相关的分区。比如查 2024 年的数据,如果 WHERE 条件里写了timestamp >= '2024-01-01' AND timestamp < '2025-01-01',数据库就只会扫daily_bars_2024这一个分区。
提示:
EXPLAIN ANALYZE是你最好的朋友。每次写完复杂查询,我都会跑一遍看执行计划,确认没有全表扫描、没有嵌套循环过多的情况。有一次一个查询跑了 8 秒,加了分区裁剪条件后降到 50 毫秒。
4. 接口层设计与常见问题排查
4.1 REST 接口的字段设计与分页策略
接口设计遵循两个原则:字段名语义化和分页必须带总数。字段名不用缩写,close_price就比cp好,虽然多打几个字符,但可读性提升巨大。分页返回里必须包含total、page、page_size三个字段,前端才能正确渲染分页控件。
from fastapi import FastAPI, Query from pydantic import BaseModel from typing import List, Optional app = FastAPI() class BarResponse(BaseModel): symbol: str timestamp: str open: float high: float low: float close: float volume: int @app.get("/api/v1/bars/daily", response_model=List[BarResponse]) async def get_daily_bars( symbol: str, start_date: str, end_date: str, page: int = Query(1, ge=1), page_size: int = Query(100, ge=1, le=1000) ): offset = (page - 1) * page_size # 查询逻辑省略 return resultspage_size的上限设为 1000,是为了防止单次请求拉取过多数据导致内存溢出。这个值可以根据实际部署环境调整,但一定要有上限。
4.2 实时推送的 WebSocket 实现要点
实时行情用 WebSocket 推送,核心要解决三个问题:连接管理、消息格式、断线重连。连接管理用一个字典维护symbol -> set(connections)的映射,有新数据时只推给订阅了该标的的连接。消息格式用 JSON,包含type、symbol、data、timestamp四个字段。
from fastapi import WebSocket from collections import defaultdict subscriptions = defaultdict(set) @app.websocket("/ws/market") async def market_ws(websocket: WebSocket): await websocket.accept() try: while True: msg = await websocket.receive_json() if msg["action"] == "subscribe": subscriptions[msg["symbol"]].add(websocket) elif msg["action"] == "unsubscribe": subscriptions[msg["symbol"]].discard(websocket) except Exception: for conns in subscriptions.values(): conns.discard(websocket)断线重连在客户端做,服务端只负责在连接断开时清理订阅关系。客户端重连后重新发送订阅请求即可。这里有个坑:如果服务端不主动清理断开的连接,内存会慢慢泄漏,所以except块里的清理逻辑是必须的。
4.3 常见问题速查表与避坑经验
| 问题现象 | 可能原因 | 排查方法 | 解决方案 |
|---|---|---|---|
| 数据拉取返回空 | 数据源接口变更或限流 | 打印原始响应内容 | 检查接口文档,调整请求参数 |
| 写入速度突然变慢 | 索引过多或分区未命中 | 查看 pg_stat_activity | 减少索引,检查 WHERE 条件 |
| 查询结果时间戳错乱 | 时区未统一 | 检查数据库时区设置 | 统一用 UTC 存储,展示时转换 |
| WebSocket 频繁断开 | 心跳缺失或代理超时 | 查看客户端和服务端日志 | 加心跳机制,调整超时时间 |
| 复权价格对不上 | 复权因子计算错误 | 手动核对除权日数据 | 重新计算复权因子表 |
我踩过最深的坑是时区问题。一开始数据库存的是本地时间,接口返回也是本地时间,但数据源给的是 UTC 时间,导致 K 线时间对不上。后来统一改成“存储用 UTC,接口返回 ISO 8601 格式带时区标识”,问题才彻底解决。这个教训是:时间戳永远不要用无时区的字符串,一定要带时区信息。
另一个坑是浮点数精度。价格用FLOAT存储会出现0.1 + 0.2 != 0.3的问题,后来全部改成NUMERIC(12,4),精度问题消失。金融数据对精度要求高,浮点数能不用就不用。
5. 部署与运维的实战建议
5.1 容器化部署与资源限制
整套服务用 Docker Compose 编排,包含 PostgreSQL、Redis、FastAPI 应用三个容器。资源限制一定要设,否则某个容器内存泄漏会把整台机器拖垮。我的配置是:PostgreSQL 限 2GB 内存,Redis 限 512MB,应用容器限 1GB。
services: db: image: timescale/timescaledb:latest-pg15 deploy: resources: limits: memory: 2G volumes: - pgdata:/var/lib/postgresql/data redis: image: redis:7-alpine deploy: resources: limits: memory: 512M app: build: . deploy: resources: limits: memory: 1G depends_on: - db - redis数据卷pgdata必须挂载到宿主机,否则容器重建数据就没了。这个坑我踩过一次,重建容器后发现所有历史数据消失,只能重新拉,花了整整一天。
5.2 监控指标与告警设置
监控三个核心指标:采集成功率、写入延迟、接口响应时间。采集成功率低于 95% 触发告警,写入延迟超过 5 秒触发告警,接口 P99 响应时间超过 1 秒触发告警。监控用 Prometheus + Grafana,应用层暴露/metrics接口。
告警渠道我用的是邮件加 webhook,没有用太复杂的方案。关键是要设置告警抑制,否则数据源故障时会在短时间内发出几百条告警,反而掩盖了真正的问题。我的做法是:同一类型的告警 5 分钟内只发一次。
5.3 数据备份与恢复演练
备份策略是每日全量 + 实时 WAL 归档。全量备份用pg_dump,每天凌晨 3 点执行,保留最近 30 天。WAL 归档让数据库可以恢复到任意时间点,对金融数据来说这个能力很重要——万一某次清洗逻辑写错了,可以精确恢复到出错之前的状态。
恢复演练我每季度做一次,流程是:从备份文件恢复到一个临时数据库,跑一遍数据完整性校验(检查记录数、检查关键字段非空、检查时间范围连续性),确认无误后才算演练通过。没做过恢复演练的备份等于没有备份,这句话在金融数据场景下尤其正确。
提示:备份文件不要和数据库放在同一台机器上。我用的是对象存储,备份完成后自动上传,本地只保留最近 3 天的文件。这样即使机器完全损坏,数据也能找回来。
5.4 后续扩展方向与个人体会
这套服务目前支撑了我自己的策略回测和几个朋友的行情看板需求,运行了大半年,整体稳定。后续我打算加两个东西:一是因子计算模块,把常用的技术指标(MA、MACD、RSI)在数据库层预计算好,查询时直接取;二是数据质量报告,每天自动生成一份数据完整性报告,推送到邮箱。
我个人在实际操作中的体会是:金融数据服务最难的不是技术,而是对数据本身的理解。什么算异常值、停牌怎么处理、复权因子怎么算,这些问题的答案不在代码里,在业务逻辑里。技术只是工具,把业务逻辑想清楚才是根本。另外,不要追求一步到位,我第一版只支持日线数据,跑通了才加分钟线,加了分钟线才加实时推送。每次只加一个功能,出问题容易定位,改起来也快。