ARTICLE DETAIL

资讯详情

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

从零搭建金融数据服务:采集、存储、计算与接口全流程

从零搭建金融数据服务:采集、存储、计算与接口全流程 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()转成标准格式最后交给存储层。统一格式我定义了一个 dataclassfrom 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_retries5, base_delay1): 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第{attempt1}次失败{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(hostlocalhost, port6379, db0) def get_quote(symbol): cache_key fquote:{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(fquote:{symbol}:latest)3.3 技术指标计算的实现与优化技术指标计算是这套服务的核心价值之一。我实现了 MA、EMA、MACD、RSI、布林带这几个常用的。以 MA 为例用 pandas 的 rolling 方法一行就能算出来import pandas as pd def calculate_ma(df, window): df[fma{window}] df[close].rolling(windowwindow).mean() return df但实际用的时候有几个坑。第一数据不足 window 条时结果是 NaN接口返回时要处理。第二计算要基于复权后的价格否则除权除息那天指标会跳变。第三增量计算比全量计算复杂得多我目前是每次全量重算最近 250 个交易日的数据简单可靠性能也扛得住。def calculate_all_indicators(symbol): df load_recent_data(symbol, days250) 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/searchGET搜索标的参数校验用 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, ge1, le1000) ): data query_history(symbol, start, end, limit) return {symbol: symbol, data: data}ge1, le1000限制了 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 从最初的一个脚本慢慢长成了现在这个有采集、存储、计算、接口四层的完整服务。中间踩的坑不少但每解决一个问题对金融数据管道的理解就深一层。如果你也在搭类似的东西建议先从最小可用版本开始别一上来就追求大而全跑通了再逐步加功能。
返回列表