尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

从零搭建金融数据聚合服务:架构设计、缓存策略与异常处理实战

发布时间:2026/9/26 5:09:25

资讯中心
01
ARTICLE

从零搭建金融数据聚合服务:架构设计、缓存策略与异常处理实战

从零搭建金融数据聚合服务:架构设计、缓存策略与异常处理实战
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或者GoNode.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/closeOHLC、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, window20, threshold3.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事件里启动。每个数据源对应一个jobjob的触发间隔和Adapter里定义的更新频率一致。错误重试这块我做了三层保护。第一层是单次请求重试在Adapter的fetch_raw里用tenacity库做重试最多重试3次每次间隔指数退避。第二层是任务级重试如果整个job执行失败APScheduler会记录错误我在job外面包了一层try-except失败后把任务状态写到数据库下一个周期再试。第三层是熔断机制如果某个数据源连续失败超过10次自动暂停该数据源的调度并发送告警通知。from tenacity import retry, stop_after_attempt, wait_exponential retry(stopstop_after_attempt(3), waitwait_exponential(multiplier1, min2, max30)) 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.04Redis 7.0PostgreSQL 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:passlocalhost: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(fValidation failed for {symbol}) return quotes filter_outliers(quotes) for q in quotes: key frt:{self.adapter.source}:{q.symbol} self.cache.set(key, q, ttlself.adapter.ttl) logger.info(fUpdated {len(quotes)} quotes for {symbol})这里有个设计决策写缓存用set而不是pipeline。因为每个quote的key不同用pipeline批量写虽然快但一旦中间出错不好排查。而且我的数据量不大单个set的性能完全够用。如果以后数据量上来了再改成pipeline也不迟。3.4 API接口设计与限流鉴权对外API我设计了三个端点/api/v1/quote/{symbol}查单个品种最新报价/api/v1/quotes?symbolsa,b,c批量查询/api/v1/history/{symbol}?startend查历史数据。限流用的是Redis的滑动窗口算法每个API key每分钟最多60次请求。鉴权用简单的Bearer Tokentoken存在数据库里每个token关联一个用户和配额。from fastapi import FastAPI, Depends, HTTPException, Header import time app FastAPI() async def rate_limit(authorization: str Header(...)): token authorization.replace(Bearer , ) key fratelimit:{token}:{int(time.time() // 60)} count redis.incr(key) if count 1: redis.expire(key, 60) if count 60: raise HTTPException(status_code429, detailRate limit exceeded) return token app.get(/api/v1/quote/{symbol}) async def get_quote(symbol: str, token: str Depends(rate_limit)): data cache.get(frt:*:{symbol}) if not data: raise HTTPException(status_code404, detailSymbol 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(fPrimary failed for {symbol}, trying fallback) try: data await fallback_adapter.fetch(symbol) return data except Exception: cached cache.get(frt:*:{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, ttlbase_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.ConnectionPoolPostgreSQL用psycopg2.pool.ThreadedConnectionPool。连接池大小根据并发量调我设的是Redis 20个连接PostgreSQL 10个连接。第二个技巧是异步化。FastAPI的接口用async defAdapter的fetch用httpx.AsyncClient这样多个请求可以并发处理不会互相阻塞。但要注意CPU密集型的操作比如数据清洗不要放在async函数里会阻塞事件循环应该用run_in_executor放到线程池里跑。第三个技巧是日志分级。DEBUG级别的日志在生产环境一定要关掉不然IO开销很大。我用的是Python的logging模块生产环境设INFO级别只记录关键操作和错误。日志格式用JSON方便后续用ELK或者Loki做聚合分析。这套financial-services从第一版到现在稳定运行了大概八个月中间经历过两次上游封禁、一次Redis内存告警、一次数据库连接池耗尽但都通过上面说的这些机制扛过来了。我个人在实际操作中的体会是金融数据服务最核心的不是技术多先进而是对异常的容忍度和恢复速度。上游一定会挂网络一定会抖数据一定会有脏的关键是在这些情况下你的服务还能不能给用户一个可用的结果。哪怕返回的是稍微旧一点的数据也比直接报错强。
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

◈

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

◐

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

▲

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。