☰
从零搭建金融数据聚合服务:架构设计、缓存策略与异常处理实战
2026/9/26 5:09:21 网站建设 项目流程

1. 金融数据服务从零搭建的完整思路

1.1 为什么我要自己搭一套金融数据服务

最早接触金融数据这块,是因为我需要一套能稳定跑在自己服务器上的行情聚合接口。市面上的商业数据API要么按调用次数收费贵得离谱,要么延迟高得让人抓狂,要么就是文档写得跟天书一样。我试过直接用某平台的免费接口,结果某天早上打开电脑发现接口挂了,没有任何预警,那感觉就像早上起来发现楼下早餐店突然关门了一样难受。

所以后来我决定自己搭一套。核心诉求其实就三个:数据要稳、延迟要低、成本要可控。这套服务我给它起名叫 financial-services,本质上是一个金融数据聚合与分发层,把多个数据源的数据拉过来,做清洗、标准化、缓存,然后通过统一的RESTful接口对外提供服务。它解决的核心问题是:让上层应用不用关心底层数据源是谁、格式什么样、什么时候会挂,只需要调我的接口就行。

这套东西适合谁呢?如果你是个独立开发者,想做个量化回测工具、盯盘小助手、或者个人记账应用里需要实时汇率,那这套方案非常适合你。如果你是小团队的技术负责人,需要给内部系统提供统一的金融数据出口,也可以直接参考。但如果你需要的是毫秒级的高频交易数据,那这套方案可能不太够用,得走专线或者托管机房,那是另一个话题了。

1.2 整体架构是怎么设计的

架构设计这块我改了三版才定下来。第一版是简单的请求转发,用户调我的接口,我实时去调上游,结果上游一慢我就跟着慢,上游一挂我就跟着挂。第二版加了本地缓存,但缓存策略太粗暴,所有数据统一五分钟过期,导致汇率这种变化快的品种数据严重滞后。第三版也就是现在这版,才算是找到了一个比较平衡的方案。

整体分四层:接入层、聚合层、缓存层、调度层。接入层负责接收外部请求和做限流鉴权;聚合层负责调用上游数据源并做数据清洗和格式统一;缓存层用Redis做多级缓存,不同品种设置不同的过期时间;调度层负责定时任务,主动去拉取数据更新缓存,而不是等用户请求来了才去拉。

为什么这么设计?核心思路是把被动变成主动。用户请求来的时候,我直接从缓存返回,响应时间稳定在10毫秒以内。后台调度器按照不同频率去更新不同品种的数据,比如外汇汇率每30秒更新一次,股票日线数据每天收盘后更新一次,加密货币因为7x24小时交易所以每10秒更新一次。这样既保证了数据的新鲜度,又不会被上游的响应速度拖累。

注意:调度频率不是越高越好。我一开始把加密货币设成每3秒更新,结果上游直接给我限流了,账号被封了24小时。后来改成10秒,一直很稳。

1.3 技术选型背后的取舍逻辑

技术栈这块我选的是Python + FastAPI + Redis + PostgreSQL + APScheduler。为什么用Python?因为金融数据处理生态里Python的库最全,pandas、numpy做数据清洗太方便了,而且FastAPI的性能在Python框架里算是第一梯队的,异步支持也好。

为什么不用Node.js或者Go?Node.js做IO密集型确实强,但数据处理这块生态不如Python成熟。Go的性能和并发确实好,但开发效率对我来说不如Python,而且很多金融数据处理的库在Go里要么没有要么不成熟。这是一个典型的开发效率 vs 运行效率的取舍,我选择了开发效率,因为这套服务的瓶颈不在计算,而在网络IO和上游限流。

数据库选PostgreSQL是因为金融数据有很多结构化查询需求,比如按时间范围查、按品种聚合、做同比环比计算,这些用关系型数据库最顺手。Redis做缓存不用多说,关键是它的过期策略和数据结构非常适合做多品种多周期的缓存管理。

APScheduler做调度是因为它够轻量,支持cron表达式和间隔触发,而且可以持久化任务状态。没用Celery是因为Celery对我来说太重了,这套服务不需要分布式任务队列那么复杂的编排。

2. 核心模块拆解与关键细节

2.1 数据源适配层怎么做到可插拔

数据源适配层是整个服务里最脏最累的活。每个上游数据源的接口格式、认证方式、返回结构、错误码都不一样,如果不做抽象,每接一个新源就要改一遍业务代码,那维护成本会爆炸。

我的做法是定义一个BaseAdapter抽象类,所有数据源适配器都继承它,必须实现三个方法:fetch_raw()负责调上游拿原始数据,normalize()负责把原始数据转成统一格式,validate()负责校验数据合理性。统一格式我定义了一个标准结构,包含symbol(品种代码)、timestamp(时间戳)、open/high/low/close(OHLC)、volume(成交量)、source(数据来源)这几个字段。

from abc import ABC, abstractmethod from dataclasses import dataclass from typing import List @dataclass class StandardQuote: symbol: str timestamp: int open: float high: float low: float close: float volume: float source: str class BaseAdapter(ABC): @abstractmethod async def fetch_raw(self, symbol: str) -> dict: pass @abstractmethod def normalize(self, raw: dict) -> List[StandardQuote]: pass @abstractmethod def validate(self, quotes: List[StandardQuote]) -> bool: pass

这样设计的好处是,接新数据源只需要写一个新的Adapter类,注册到工厂里就行,业务代码一行不用改。我目前接了四个源:一个外汇数据源、一个加密货币数据源、一个股票数据源、一个贵金属数据源。每个源的更新频率和限流策略都在Adapter里自己管理。

实操心得:写Adapter的时候一定要处理上游返回空数据的情况。我有一次遇到上游返回了200状态码但body是空的,结果normalize直接抛异常,整个调度任务挂了。后来在fetch_raw里加了空值检查,返回空列表而不是抛异常,调度器就能继续跑下一个任务。

2.2 数据清洗与异常值过滤的实操方法

上游数据不是拿来就能用的,里面有很多脏数据。我遇到过的情况包括:价格突然变成0、成交量是负数、时间戳是未来时间、同一个时间点有两条不同价格的数据。这些如果不处理,上层应用拿到的数据就是垃圾。

我的清洗流程分三步。第一步是基础校验,检查价格是否大于0、成交量是否非负、时间戳是否在合理范围内(比如不能超过当前时间5分钟)。第二步是异常值检测,我用的是基于滑动窗口的Z-score方法,如果某个价格点偏离过去20个点的均值超过3个标准差,就标记为异常。第三步是去重,同一个symbol同一个timestamp只保留最新的一条。

import numpy as np def filter_outliers(quotes, window=20, threshold=3.0): if len(quotes) < window: return quotes closes = np.array([q.close for q in quotes]) filtered = [] for i in range(len(quotes)): if i < window: filtered.append(quotes[i]) continue window_data = closes[i-window:i] mean = np.mean(window_data) std = np.std(window_data) if std == 0: filtered.append(quotes[i]) continue z_score = abs(closes[i] - mean) / std if z_score < threshold: filtered.append(quotes[i]) return filtered

这里有个细节要注意:Z-score的窗口大小和阈值需要根据品种调整。外汇波动小,window设20、threshold设3就够了。加密货币波动大,window得设50、threshold得设5,不然正常的大波动会被误杀。这个参数没有万能值,得根据实际数据调。

2.3 多级缓存策略与过期时间设计

缓存这块我踩的坑最多。最开始所有数据统一5分钟过期,结果用户查实时汇率,拿到的是5分钟前的价格,投诉说数据不准。后来改成所有数据都实时拉,结果上游限流把我封了。

现在的策略是按品种类型分级。我在Redis里用不同的key前缀区分数据类型,每种类型设置不同的TTL。具体配置如下表:

数据类型Redis Key前缀TTL更新方式
实时汇率rt:fx:30秒调度器主动更新
加密货币rt:crypto:10秒调度器主动更新
股票实时rt:stock:60秒调度器主动更新
股票日线daily:stock:24小时收盘后更新
贵金属rt:metal:60秒调度器主动更新
历史数据hist:7天按需拉取后缓存

为什么这么设?实时汇率30秒是因为外汇市场虽然24小时交易,但主要波动集中在特定时段,30秒的延迟对大多数应用够用了。加密货币10秒是因为7x24小时交易且波动剧烈,用户对延迟更敏感。股票实时60秒是因为A股本身3秒一个tick,但我的上游数据源最快也就1分钟更新一次,设更短没意义。

注意:TTL不要设得太短,否则缓存还没被用到就过期了,等于白缓存。也不要设太长,否则数据陈旧。我的经验是TTL设为上游更新频率的1.5到2倍比较合适。

2.4 调度器的任务编排与错误重试

调度器用的是APScheduler的BackgroundScheduler,在FastAPI的startup事件里启动。每个数据源对应一个job,job的触发间隔和Adapter里定义的更新频率一致。

错误重试这块我做了三层保护。第一层是单次请求重试,在Adapter的fetch_raw里用tenacity库做重试,最多重试3次,每次间隔指数退避。第二层是任务级重试,如果整个job执行失败,APScheduler会记录错误,我在job外面包了一层try-except,失败后把任务状态写到数据库,下一个周期再试。第三层是熔断机制,如果某个数据源连续失败超过10次,自动暂停该数据源的调度,并发送告警通知。

from tenacity import retry, stop_after_attempt, wait_exponential @retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=2, max=30)) async def fetch_with_retry(adapter, symbol): return await adapter.fetch_raw(symbol)

熔断这块我用了一个简单的计数器,存在Redis里,key是circuit:{source_name},每次失败INCR,成功就DEL。当计数超过阈值时,调度器跳过该数据源的job,并记录日志。等过了冷却期(我设的是30分钟),自动恢复。

3. 完整实操流程与核心代码实现

3.1 环境准备与依赖安装

先说环境。我用的Python 3.11,太老的版本不支持一些新语法,太新的版本有些库还没适配。操作系统是Ubuntu 22.04,Redis 7.0,PostgreSQL 15。这些版本都是我实测稳定的组合。

依赖安装用pip就行,主要依赖如下:

pip install fastapi uvicorn redis psycopg2-binary apscheduler tenacity numpy pandas httpx pydantic

这里重点说几个库的选择理由。HTTP客户端用httpx而不是requests,因为httpx支持异步,在FastAPI的异步上下文里不会阻塞事件循环。数据库驱动用psycopg2-binary而不是asyncpg,因为我的数据库操作不频繁,同步驱动够用且更稳定。pydantic用来做数据校验和序列化,FastAPI原生支持,非常顺手。

实操心得:psycopg2-binary在有些系统上安装会报错,需要先装libpq-dev。如果遇到编译错误,直接apt install libpq-dev再pip install就好了。

3.2 项目目录结构与配置管理

目录结构我参考了FastAPI官方推荐的项目布局,但做了一些调整:

financial-services/ ├── app/ │ ├── main.py │ ├── config.py │ ├── adapters/ │ │ ├── base.py │ │ ├── fx_adapter.py │ │ ├── crypto_adapter.py │ │ └── stock_adapter.py │ ├── services/ │ │ ├── aggregator.py │ │ ├── cache.py │ │ └── cleaner.py │ ├── scheduler/ │ │ └── jobs.py │ ├── api/ │ │ └── routes.py │ └── models/ │ └── schemas.py ├── tests/ ├── requirements.txt └── .env

配置管理用pydantic的BaseSettings,从环境变量和.env文件读取配置。这样本地开发和线上部署可以用同一套代码,只是环境变量不同。

from pydantic_settings import BaseSettings class Settings(BaseSettings): redis_url: str = "redis://localhost:6379/0" database_url: str = "postgresql://user:pass@localhost:5432/finance" fx_api_key: str = "" crypto_api_key: str = "" log_level: str = "INFO" class Config: env_file = ".env" settings = Settings()

3.3 核心聚合服务的实现细节

聚合服务是连接Adapter和缓存的桥梁。它的工作流程是:调度器触发 -> 聚合服务调用Adapter的fetch_raw -> normalize -> validate -> filter_outliers -> 写入Redis -> 写入PostgreSQL(可选)。

class AggregatorService: def __init__(self, adapter: BaseAdapter, cache: CacheService): self.adapter = adapter self.cache = cache async def update(self, symbol: str): raw = await fetch_with_retry(self.adapter, symbol) quotes = self.adapter.normalize(raw) if not self.adapter.validate(quotes): logger.warning(f"Validation failed for {symbol}") return quotes = filter_outliers(quotes) for q in quotes: key = f"rt:{self.adapter.source}:{q.symbol}" self.cache.set(key, q, ttl=self.adapter.ttl) logger.info(f"Updated {len(quotes)} quotes for {symbol}")

这里有个设计决策:写缓存用set而不是pipeline。因为每个quote的key不同,用pipeline批量写虽然快,但一旦中间出错不好排查。而且我的数据量不大,单个set的性能完全够用。如果以后数据量上来了,再改成pipeline也不迟。

3.4 API接口设计与限流鉴权

对外API我设计了三个端点:/api/v1/quote/{symbol}查单个品种最新报价,/api/v1/quotes?symbols=a,b,c批量查询,/api/v1/history/{symbol}?start=&end=查历史数据。

限流用的是Redis的滑动窗口算法,每个API key每分钟最多60次请求。鉴权用简单的Bearer Token,token存在数据库里,每个token关联一个用户和配额。

from fastapi import FastAPI, Depends, HTTPException, Header import time app = FastAPI() async def rate_limit(authorization: str = Header(...)): token = authorization.replace("Bearer ", "") key = f"ratelimit:{token}:{int(time.time() // 60)}" count = redis.incr(key) if count == 1: redis.expire(key, 60) if count > 60: raise HTTPException(status_code=429, detail="Rate limit exceeded") return token @app.get("/api/v1/quote/{symbol}") async def get_quote(symbol: str, token: str = Depends(rate_limit)): data = cache.get(f"rt:*:{symbol}") if not data: raise HTTPException(status_code=404, detail="Symbol not found") return data

注意:限流的key一定要带时间窗口,不然计数器永远不重置。我一开始忘了加时间戳,结果用户第一次请求就把配额用完了,后面永远429。

3.5 数据持久化与历史查询优化

实时数据放Redis,历史数据放PostgreSQL。历史数据表按symbol和时间做联合索引,查询的时候用WHERE symbol = ? AND timestamp BETWEEN ? AND ?,走索引速度很快。

CREATE TABLE quotes ( id BIGSERIAL PRIMARY KEY, symbol VARCHAR(20) NOT NULL, timestamp BIGINT NOT NULL, open NUMERIC(18,8), high NUMERIC(18,8), low NUMERIC(18,8), close NUMERIC(18,8), volume NUMERIC(24,8), source VARCHAR(20), created_at TIMESTAMP DEFAULT NOW() ); CREATE INDEX idx_quotes_symbol_ts ON quotes (symbol, timestamp DESC);

历史数据的写入我用的是批量插入,每1000条一批,用execute_values比逐条插入快几十倍。但要注意,批量插入的时候如果有一条数据格式不对,整批都会失败。所以我在插入前会先做一次全量校验,确保数据干净。

4. 常见问题与排查技巧实录

4.1 上游接口突然不可用怎么办

这是最常见的问题。上游接口挂掉的原因五花八门:服务器维护、限流封禁、接口改版、网络抖动。我的处理策略是分级降级。

第一级:如果只是单次请求失败,tenacity自动重试,用户无感知。第二级:如果连续失败超过5次,切换到备用数据源(如果有的话)。第三级:如果没有备用源,返回缓存中的最后一条数据,并在响应头里加X-Data-Stale: true标记。第四级:如果连缓存都没有,返回503并附带预计恢复时间。

async def get_quote_with_fallback(symbol: str): try: data = await primary_adapter.fetch(symbol) return data except Exception: logger.warning(f"Primary failed for {symbol}, trying fallback") try: data = await fallback_adapter.fetch(symbol) return data except Exception: cached = cache.get(f"rt:*:{symbol}") if cached: cached["stale"] = True return cached raise HTTPException(503, "Service temporarily unavailable")

实操心得:备用数据源不一定要跟主源一模一样,可以是不同粒度的。比如主源提供分钟线,备用源只提供小时线,那降级的时候至少还能用,虽然精度差了但总比没有强。

4.2 数据延迟突然变大的排查思路

数据延迟变大通常有三个原因:上游变慢、调度器卡住、Redis变慢。排查顺序是从外到内。

先看上游的响应时间,我在Adapter里记录了每次请求的耗时,写到日志里。如果上游耗时从200ms涨到2s,那就是上游的问题,只能等或者切备用源。如果上游正常,就看调度器的执行日志,看是不是某个job执行时间过长阻塞了其他job。APScheduler默认是单线程执行job的,如果一个job卡住,后面的都会排队。解决办法是给scheduler配置线程池,或者把耗时的job改成异步执行。

如果前两个都正常,那就是Redis的问题。用redis-cli --latency看延迟,如果超过1ms就要注意了。常见原因是Redis内存满了触发淘汰,或者有大key导致阻塞。我遇到过一次是因为某个key存了一个巨大的列表,每次读取都要几毫秒,后来把大key拆成多个小key就好了。

4.3 缓存穿透和缓存雪崩的预防

缓存穿透是指查一个不存在的symbol,每次都打到数据库。我的预防措施是空值缓存:如果查数据库没查到,就在Redis里存一个空标记,TTL设短一点比如60秒。这样下次查同一个不存在的symbol,直接返回空,不会打到数据库。

缓存雪崩是指大量key同时过期,请求全部打到上游。我的预防措施是TTL加随机抖动:在基础TTL上加上一个0到10秒的随机值,这样key不会在同一秒集中过期。

import random def set_with_jitter(key, value, base_ttl): jitter = random.randint(0, 10) cache.set(key, value, ttl=base_ttl + jitter)

4.4 常见问题速查表

问题现象可能原因排查方法解决方案
接口返回429触发限流看Redis计数器降低调度频率或申请更高配额
数据价格明显错误上游脏数据对比多个数据源加强validate和filter_outliers
调度任务不执行scheduler挂了看APScheduler日志重启服务,检查线程池配置
Redis内存暴涨大key或TTL过长redis-cli --bigkeys拆分大key,缩短TTL
数据库查询慢缺索引或数据量大EXPLAIN ANALYZE加索引,考虑分区表
服务启动报错依赖缺失或配置错误看启动日志检查.env和依赖版本

4.5 性能优化的几个实用技巧

第一个技巧是连接池。Redis和PostgreSQL都要用连接池,不要每次请求都新建连接。Redis用redis.ConnectionPool,PostgreSQL用psycopg2.pool.ThreadedConnectionPool。连接池大小根据并发量调,我设的是Redis 20个连接,PostgreSQL 10个连接。

第二个技巧是异步化。FastAPI的接口用async def,Adapter的fetch用httpx.AsyncClient,这样多个请求可以并发处理,不会互相阻塞。但要注意,CPU密集型的操作比如数据清洗不要放在async函数里,会阻塞事件循环,应该用run_in_executor放到线程池里跑。

第三个技巧是日志分级。DEBUG级别的日志在生产环境一定要关掉,不然IO开销很大。我用的是Python的logging模块,生产环境设INFO级别,只记录关键操作和错误。日志格式用JSON,方便后续用ELK或者Loki做聚合分析。

这套financial-services从第一版到现在稳定运行了大概八个月,中间经历过两次上游封禁、一次Redis内存告警、一次数据库连接池耗尽,但都通过上面说的这些机制扛过来了。我个人在实际操作中的体会是,金融数据服务最核心的不是技术多先进,而是对异常的容忍度和恢复速度。上游一定会挂,网络一定会抖,数据一定会有脏的,关键是在这些情况下你的服务还能不能给用户一个可用的结果。哪怕返回的是稍微旧一点的数据,也比直接报错强。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询