1. 金融数据服务从零搭建的完整思路
1.1 为什么我要自己动手做一套金融数据服务
最早接触金融数据这块,是因为我需要一套能稳定拉取行情、做基础指标计算、再对外提供查询接口的服务。市面上现成的方案要么太贵,要么数据延迟高得离谱,要么接口限制多得让人抓狂。折腾了一圈之后我决定自己搭一套,项目代号就叫financial-services。
这套东西说白了就是一个中间层:上游对接公开的行情数据源,中间做清洗、存储、计算,下游通过 REST 接口对外提供标准化的金融数据。它能解决的核心问题有三个——数据源格式不统一、实时计算逻辑重复造轮子、查询接口各自为政。适合谁参考呢?有一定后端基础、想了解金融数据管道怎么搭的开发者,或者手里有小规模量化策略、需要自己维护数据链路的个人交易者。
我先把整体架构讲清楚,再逐层拆解每个模块的设计取舍和踩坑记录。整套服务我用 Python 写的,Web 框架选的 FastAPI,数据库用 PostgreSQL 加 Redis 做缓存,定时任务用 APScheduler。选型理由后面会细说,先把骨架立起来。
1.2 整体架构分层与数据流向
整个服务分成四层,从下往上依次是:数据采集层、数据存储层、计算服务层、接口暴露层。数据采集层负责从外部拉原始数据,做初步的格式归一化;存储层把清洗后的数据落到 PostgreSQL,热点数据丢进 Redis;计算层跑各种技术指标和统计聚合;接口层用 FastAPI 暴露 RESTful 端点。
数据流向是这样的:定时任务触发采集器,采集器把原始 JSON 写进一个临时表,清洗模块读取临时表做字段映射和异常值过滤,写入正式表。同时,计算模块监听正式表的数据变更,触发指标重算,结果写入指标表并刷新 Redis 缓存。接口层查询时优先读 Redis,未命中再查 PostgreSQL,查完回写缓存。
这个分层的好处是每层职责单一,出问题好定位。比如某天发现接口返回的 MA5 不对,我可以先查 Redis 缓存是不是脏了,再查计算层是不是没触发,最后查采集层数据是不是缺了。如果全揉在一起,排查起来就是一团乱麻。
1.3 技术选型的取舍逻辑
选 Python 而不是 Java 或 Go,主要考虑的是开发效率和生态。金融数据处理这块,pandas、numpy 这些库太成熟了,用 Java 写同样的逻辑代码量至少翻倍。FastAPI 相比 Flask 和 Django,异步性能好,自带 OpenAPI 文档,对于要对外提供接口的场景非常合适。
数据库选 PostgreSQL 而不是 MySQL,是因为 PostgreSQL 对 JSON 字段的支持更好,而且窗口函数、CTE 这些高级特性在处理金融数据时特别顺手。Redis 用来做缓存和简单的发布订阅,不引入 Kafka 是因为项目规模还没到那个量级,用 Redis 的 Pub/Sub 足够。
注意:如果你的数据量级到了每天千万条以上,PostgreSQL 单表会扛不住,需要考虑 TimescaleDB 或者 ClickHouse 这类时序数据库。我当前的数据量在每天几十万条,PostgreSQL 完全够用。
2. 数据采集层的核心细节与实操要点
2.1 数据源对接的通用抽象设计
金融数据源五花八门,有返回 JSON 的,有返回 CSV 的,还有需要 WebSocket 长连接的。为了不让采集逻辑散落各处,我定义了一个抽象基类BaseCollector,所有具体采集器都继承它,实现fetch()和parse()两个方法。
from abc import ABC, abstractmethod class BaseCollector(ABC): @abstractmethod def fetch(self, **kwargs): """从数据源拉取原始数据""" pass @abstractmethod def parse(self, raw_data): """将原始数据解析为统一格式""" pass这样做的好处是,新增一个数据源只需要写一个新的子类,不用动调度逻辑。调度器统一调用fetch()拿原始数据,再调parse()转成标准格式,最后交给存储层。
统一格式我定义了一个 dataclass:
from dataclasses import dataclass from datetime import datetime @dataclass class MarketData: symbol: str timestamp: datetime open: float high: float low: float close: float volume: float source: str所有采集器最终都产出这个结构,下游就不用关心数据是从哪来的了。
2.2 采集频率与重试机制的设计
采集频率不能拍脑袋定。日线数据每天收盘后拉一次就行,分钟线数据要看你的策略需求,但频率越高对数据源的压力越大,被封的风险也越高。我目前是日线每天下午四点拉一次,分钟线每五分钟拉一次。
重试机制是必须的,网络抖动、数据源限流都会导致单次请求失败。我用的是指数退避策略:第一次失败等 1 秒,第二次等 2 秒,第三次等 4 秒,最多重试 5 次。超过 5 次就记录错误日志并跳过,等下一轮调度再补。
import time import logging def retry_with_backoff(func, max_retries=5, base_delay=1): for attempt in range(max_retries): try: return func() except Exception as e: if attempt == max_retries - 1: logging.error(f"重试{max_retries}次后仍失败: {e}") raise delay = base_delay * (2 ** attempt) logging.warning(f"第{attempt+1}次失败,{delay}秒后重试") time.sleep(delay)实操心得:重试的时候一定要加随机抖动,比如
delay = base_delay * (2 ** attempt) + random.uniform(0, 1)。否则多个采集任务同时失败后会在同一时刻重试,形成惊群效应,把数据源直接打挂。
2.3 数据清洗的边界条件处理
原始数据里什么妖魔鬼怪都有:价格是字符串的、时间戳是毫秒的、成交量是负数的、字段名大小写不一致的。清洗模块要做的就是把这些统一成标准格式。
我列了一个清洗检查清单,每次新增数据源都对照着过一遍:
| 检查项 | 处理方式 | 示例 |
|---|---|---|
| 价格字段类型 | 强制转 float,失败则丢弃 | "12.5" → 12.5 |
| 时间戳单位 | 统一转秒级 | 1700000000000 → 1700000000 |
| 字段名规范 | 统一转小写 | "Open" → "open" |
| 异常值过滤 | 价格<=0 或 成交量<0 丢弃 | -1.5 → 丢弃 |
| 缺失值处理 | 关键字段缺失丢弃,非关键填默认值 | volume 缺失填 0 |
清洗逻辑我单独写了一个模块,不跟采集混在一起。这样调试的时候可以拿历史原始数据反复跑清洗,不用重新拉数据。
3. 存储层设计与计算服务的实现
3.1 PostgreSQL 表结构设计与索引优化
正式表我按数据频率分了两种:daily_market_data和minute_market_data。字段结构一样,只是数据量和查询模式不同。
CREATE TABLE daily_market_data ( id BIGSERIAL PRIMARY KEY, symbol VARCHAR(20) NOT NULL, trade_date DATE NOT NULL, open NUMERIC(12, 4), high NUMERIC(12, 4), low NUMERIC(12, 4), close NUMERIC(12, 4), volume BIGINT, source VARCHAR(50), created_at TIMESTAMP DEFAULT NOW(), UNIQUE(symbol, trade_date) );唯一索引建在(symbol, trade_date)上,这样重复采集不会产生脏数据,用ON CONFLICT DO UPDATE就能实现幂等写入。
INSERT INTO daily_market_data (symbol, trade_date, open, high, low, close, volume, source) VALUES (%s, %s, %s, %s, %s, %s, %s, %s) ON CONFLICT (symbol, trade_date) DO UPDATE SET open = EXCLUDED.open, high = EXCLUDED.high, low = EXCLUDED.low, close = EXCLUDED.close, volume = EXCLUDED.volume;查询索引方面,除了唯一索引,我还建了一个(symbol, trade_date DESC)的复合索引,因为最常用的查询就是"取某只股票最近 N 天的数据"。
注意:NUMERIC 类型比 FLOAT 更适合存价格,因为浮点数有精度问题。0.1 + 0.2 在浮点里不等于 0.3,金融数据对精度要求高,必须用 NUMERIC。
3.2 Redis 缓存策略与失效机制
Redis 主要缓存两类数据:最新的行情快照和技术指标计算结果。缓存 key 的设计我用了{数据类型}:{标的}:{周期}的格式,比如quote:AAPL:latest、ma:AAPL:5。
失效策略用的是主动更新 + 过期兜底。数据更新时主动删除相关缓存 key,下次查询时重新加载。同时每个 key 设一个 TTL,比如行情快照 60 秒,指标数据 300 秒,防止主动删除逻辑出 bug 导致缓存永远不更新。
import redis import json r = redis.Redis(host='localhost', port=6379, db=0) def get_quote(symbol): cache_key = f"quote:{symbol}:latest" cached = r.get(cache_key) if cached: return json.loads(cached) # 缓存未命中,查数据库 data = query_from_db(symbol) r.setex(cache_key, 60, json.dumps(data)) return data def invalidate_quote(symbol): r.delete(f"quote:{symbol}:latest")3.3 技术指标计算的实现与优化
技术指标计算是这套服务的核心价值之一。我实现了 MA、EMA、MACD、RSI、布林带这几个常用的。以 MA 为例,用 pandas 的 rolling 方法一行就能算出来:
import pandas as pd def calculate_ma(df, window): df[f'ma{window}'] = df['close'].rolling(window=window).mean() return df但实际用的时候有几个坑。第一,数据不足 window 条时结果是 NaN,接口返回时要处理。第二,计算要基于复权后的价格,否则除权除息那天指标会跳变。第三,增量计算比全量计算复杂得多,我目前是每次全量重算最近 250 个交易日的数据,简单可靠,性能也扛得住。
def calculate_all_indicators(symbol): df = load_recent_data(symbol, days=250) df = calculate_ma(df, 5) df = calculate_ma(df, 10) df = calculate_ma(df, 20) df = calculate_ema(df, 12) df = calculate_ema(df, 26) df = calculate_macd(df) df = calculate_rsi(df, 14) save_indicators(symbol, df)实操心得:全量重算虽然简单,但要注意数据库写入量。250 天 × 每只股票 × 多个指标,数据量不小。我的做法是只写入最近 30 天的指标,更早的数据按需计算,不落库。
4. 接口层实现与常见问题排查
4.1 FastAPI 接口设计与参数校验
接口层我用 FastAPI 实现,主要提供这几个端点:
| 端点 | 方法 | 说明 |
|---|---|---|
/api/v1/quote/{symbol} | GET | 获取最新行情 |
/api/v1/history/{symbol} | GET | 获取历史K线 |
/api/v1/indicator/{symbol} | GET | 获取技术指标 |
/api/v1/search | GET | 搜索标的 |
参数校验用 Pydantic 模型,比如历史K线接口的查询参数:
from fastapi import FastAPI, Query from pydantic import BaseModel from datetime import date app = FastAPI() @app.get("/api/v1/history/{symbol}") async def get_history( symbol: str, start: date = Query(..., description="开始日期"), end: date = Query(..., description="结束日期"), limit: int = Query(100, ge=1, le=1000) ): data = query_history(symbol, start, end, limit) return {"symbol": symbol, "data": data}ge=1, le=1000限制了 limit 的范围,防止有人传个 100000 把数据库拖垮。
4.2 接口性能优化的三个关键点
第一个是分页。历史数据接口必须分页,不能一次返回全部。我用的是 limit + offset 的方式,简单直接。
第二个是字段裁剪。不是所有场景都需要全部字段,我加了一个fields参数,让调用方指定需要哪些字段,减少传输量。
第三个是异步查询。FastAPI 支持 async 路由,数据库查询用 asyncpg 或者把同步查询放到线程池里跑,避免阻塞事件循环。
from fastapi import FastAPI import asyncio @app.get("/api/v1/quote/{symbol}") async def get_quote(symbol: str): # 把同步的数据库查询放到线程池 loop = asyncio.get_event_loop() data = await loop.run_in_executor(None, query_quote, symbol) return data4.3 常见问题速查表
这套服务跑了大半年,遇到的问题不少,我整理了一个速查表:
| 问题现象 | 可能原因 | 排查方法 | 解决方案 |
|---|---|---|---|
| 接口返回空数据 | 采集任务失败 | 查采集日志 | 手动触发补采 |
| 指标数值异常 | 复权处理错误 | 对比原始数据 | 检查复权因子 |
| 接口响应慢 | 缓存未命中 | 查 Redis 命中率 | 调整 TTL 或预热缓存 |
| 数据重复 | 唯一索引失效 | 查表约束 | 重建唯一索引 |
| 定时任务不执行 | 调度器挂了 | 查调度日志 | 加进程守护 |
避坑技巧:定时任务一定要加监控告警。我有一次 APScheduler 因为一个未捕获的异常整个挂掉了,三天后才发现数据断了。后来加了一个心跳检测,每五分钟往 Redis 写一个时间戳,外部监控发现时间戳超过十分钟没更新就告警。
4.4 数据一致性保障的实践经验
金融数据最怕的就是不一致。同一只股票,行情接口返回的收盘价和指标接口用的收盘价对不上,那整个服务就不可信了。
我的做法是单一数据源原则:所有下游计算都从同一张正式表读数据,不允许任何模块直接访问原始数据。正式表的数据写入必须经过清洗模块,清洗模块的写入是事务性的,要么全成功要么全回滚。
def save_market_data(conn, data_list): try: with conn.cursor() as cur: for data in data_list: cur.execute(INSERT_SQL, data) conn.commit() except Exception as e: conn.rollback() raise e另外,我每天会跑一次数据校验任务,检查当天数据的完整性:有没有缺失的标的、有没有价格跳变超过 20% 的异常记录、有没有成交量为零的。发现问题就发告警,人工介入处理。
这套 financial-services 从最初的一个脚本,慢慢长成了现在这个有采集、存储、计算、接口四层的完整服务。中间踩的坑不少,但每解决一个问题,对金融数据管道的理解就深一层。如果你也在搭类似的东西,建议先从最小可用版本开始,别一上来就追求大而全,跑通了再逐步加功能。