拓冰建站拓冰建站
首页 / 资讯中心 / 正文

从零搭建金融数据服务:架构设计、核心模块与实操避坑指南

1. 金融数据服务从零搭建的完整思路1.1 为什么我要自己动手做一套金融数据服务先说清楚这个项目到底在干什么。financial-services这个名字听起来很泛实际上我做的事情是搭建一套面向个人开发者和小型团队使用的金融数据聚合与分发服务。它要解决的问题很具体——市面上现成的金融数据接口要么贵得离谱要么免费的限制多到没法用要么数据质量参差不齐。我需要一个自己能掌控的中间层把多个数据源统一收口做清洗、缓存、限流、格式化然后以统一的 API 暴露给上层应用。这套东西适合谁如果你正在做量化策略回测、个人记账工具、投资组合追踪、财经资讯聚合或者单纯想练手后端服务架构那这套思路你直接拿去改改就能用。它不涉及任何交易执行环节纯粹是数据层面的搬运和加工所以安全边界很清晰。我选择自己搭而不是直接用第三方聚合服务核心原因有三个。第一成本可控。免费额度用完之后商业接口按调用次数计费一个中等频率的策略回测跑一晚上就能烧掉几百块。第二数据主权。自己存一份历史数据想怎么查就怎么查不用担心接口突然下线或者改字段。第三定制空间。不同策略对数据频率、复权方式、时间对齐的要求完全不同只有自己掌控管道才能灵活调整。1.2 整体架构是怎么设计的架构上我采用的是经典的分层设计从上到下依次是接入层、聚合层、存储层、采集层。这个顺序跟数据流向是反的但理解起来更符合“我要什么”的思维习惯。接入层负责对外提供 API处理鉴权、限流、参数校验、响应格式化。聚合层是核心业务逻辑负责把多个数据源的原始数据做对齐、去重、补全、计算衍生指标。存储层分两块热数据放内存缓存温冷数据放时序数据库和关系库。采集层负责跟各个上游数据源打交道处理重试、断点续传、增量更新。为什么这么分因为每一层的变更频率完全不同。接入层的 API 版本可能半年才动一次聚合层的业务规则可能每周都在调采集层的上游接口可能随时变化。分层之后改任何一层都不会波及另外三层维护成本大幅降低。提示不要一上来就追求微服务。我一开始就把采集和聚合拆成两个独立服务结果调试的时候光日志追踪就烦死人。后来合并成一个进程内的模块化单体开发效率直接翻倍。等日调用量过百万再拆不迟。1.3 技术选型背后的取舍逻辑语言我选了 Python原因很直接金融数据处理生态最成熟。pandas、numpy、polars 这些库做数据清洗和计算效率比手写循环高两个数量级。虽然 Python 在超高并发场景下不如 Go 或 Rust但金融数据服务的瓶颈通常在 IO 和数据处理不在语言本身的执行速度。Web 框架用的 FastAPI。选它不是因为性能最好而是因为自动生成 OpenAPI 文档这一点太省事了。前端同事要对接的时候直接把/docs甩过去就行省掉大量沟通成本。而且 Pydantic 做参数校验和序列化类型提示写清楚之后很多低级错误在编码阶段就能发现。数据库组合是 PostgreSQL Redis Parquet 文件。PostgreSQL 存元数据、用户配置、任务调度记录这些结构化程度高的数据。Redis 做热点数据缓存和分布式锁。历史行情数据直接落 Parquet 文件按日期分区查询的时候用 DuckDB 或者 Polars 直接扫文件比塞进关系库快得多也省空间。组件选型核心考量替代方案语言Python 3.11数据生态成熟Go并发强但数据科学生态弱Web框架FastAPI自动文档类型校验Flask太裸、Django太重关系库PostgreSQLJSON支持好稳定MySQLJSON弱缓存Redis数据结构丰富Memcached功能单一列式存储Parquet压缩比高生态好CSV太慢、HDF5生态窄2. 核心模块的细节拆解与实操要点2.1 数据采集模块的容错设计采集模块是整个服务的地基它挂了后面全完。我在这个模块上踩的坑最多所以展开讲。第一个关键点是多源冗余。任何一个免费数据源都不可靠可能今天限流、明天改字段、后天直接关停。我的做法是同时接入至少三个来源每个来源写一个独立的 adapter统一实现同一个抽象基类的接口。基类定义三个方法fetch_raw()拉取原始数据、normalize()转成标准格式、health_check()自检可用性。from abc import ABC, abstractmethod from typing import List from datetime import date class BaseDataSource(ABC): abstractmethod def fetch_raw(self, symbol: str, start: date, end: date) - list: 从上游拉取原始数据 pass abstractmethod def normalize(self, raw_data: list) - List[dict]: 将原始数据转为标准格式 pass abstractmethod def health_check(self) - bool: 检查数据源是否可用 pass这样做的好处是当主源出问题时聚合层只需要切换一个配置项就能切到备用源代码零改动。我实测下来三个源同时挂掉的概率极低基本能保证 99.5% 以上的可用性。第二个关键点是增量更新与断点续传。全量拉取历史数据只在初始化时做一次之后每天只拉增量。我在 PostgreSQL 里维护一张sync_state表记录每个数据源、每个标的、每个频率的最后同步时间戳。每次采集任务启动时先读这张表从上次成功的位置继续。CREATE TABLE sync_state ( source_name VARCHAR(50) NOT NULL, symbol VARCHAR(20) NOT NULL, freq VARCHAR(10) NOT NULL, last_sync_ts TIMESTAMP NOT NULL, status VARCHAR(20) DEFAULT success, retry_count INT DEFAULT 0, PRIMARY KEY (source_name, symbol, freq) );注意时间戳一定要用 UTC 存储展示的时候再转本地时区。我一开始混用本地时间和 UTC导致跨日数据对不上排查了整整一个下午。2.2 数据清洗与标准化的具体规则原始数据拿到手之后不能直接用。不同来源的字段名、时间格式、复权方式、缺失值处理都不一样。清洗模块要做的事情就是把这些差异全部抹平。字段映射我用 YAML 配置文件管理每个数据源一个文件。这样加新源的时候不用改代码写个配置就行。# config/mappings/source_a.yaml fields: trade_date: date open_price: open high_price: high low_price: low close_price: close vol: volume amount: turnover date_format: %Y%m%d adjust: qfq # 前复权时间对齐是清洗里最麻烦的部分。不同市场、不同品种的交易时间不同有的还有午休。我的处理方式是统一转成 UTC 时间戳然后在聚合层根据标的所属市场做对齐。对于日频数据统一用交易日收盘时间作为时间戳对于分钟频数据统一用 bar 的结束时间。缺失值处理分三种情况。第一种是停牌导致的缺失这种直接标记为suspended不填充。第二种是数据源漏推这种用前后值的线性插值补上但打上interpolated标记。第三种是确实没有数据比如新上市标的的历史数据这种直接留空不硬造。异常值检测我用的是 MAD中位数绝对偏差方法比标准差更抗离群点。具体做法是计算每个标的收益率序列的 MAD超过 5 倍 MAD 的点标记为可疑人工复核后再决定是修正还是剔除。2.3 缓存策略与性能优化金融数据服务的读请求远多于写请求缓存是性能的生命线。我的缓存分三级进程内 LRU、Redis、Parquet 文件。进程内 LRU 用cachetools库缓存最近查询的 1000 个结果TTL 设 60 秒。这一级命中率大概 30%但响应时间在微秒级对高频重复查询帮助很大。Redis 缓存 TTL 设 5 分钟到 1 小时不等取决于数据频率。日频数据缓存 1 小时分钟频缓存 5 分钟。缓存 key 的设计很关键我用{freq}:{symbol}:{start}:{end}:{adjust}的格式保证不同参数组合不会串。import hashlib import json from redis import Redis def make_cache_key(freq: str, symbol: str, start: str, end: str, adjust: str) - str: raw f{freq}:{symbol}:{start}:{end}:{adjust} return ffin:data:{hashlib.md5(raw.encode()).hexdigest()}Parquet 文件是最后一道防线也是历史数据的唯一真相来源。我按{freq}/{symbol}/{year}/{month}.parquet的路径组织查询的时候用谓词下推只扫需要的分区。实测下来查单标的 5 年日线数据从 Parquet 读取比从 PostgreSQL 快 8 到 10 倍。实操心得Parquet 文件不要写太小。我一开始按天分区结果产生了几十万个小文件文件系统的 inode 都快用完了。后来改成按月分区单文件大小控制在 50MB 到 200MB 之间读写效率明显提升。3. 完整实操流程与关键环节实现3.1 环境搭建与依赖安装先把基础环境跑起来。我假设你用的是 Ubuntu 22.04 或者 macOSPython 版本 3.11 以上。# 创建虚拟环境 python3.11 -m venv venv source venv/bin/activate # 安装核心依赖 pip install fastapi uvicorn[standard] pydantic pydantic-settings pip install pandas polars pyarrow duckdb pip install sqlalchemy psycopg2-binary alembic pip install redis cachetools httpx tenacity pip install python-dateutil pytz pyyaml数据库用 Docker 起省得污染本机环境。docker run -d --name fin-pg \ -e POSTGRES_USERfin \ -e POSTGRES_PASSWORDfin_dev_pass \ -e POSTGRES_DBfinancial_services \ -p 5432:5432 \ postgres:16-alpine docker run -d --name fin-redis \ -p 6379:6379 \ redis:7-alpine目录结构我习惯这样组织清晰且容易扩展financial-services/ ├── app/ │ ├── api/ # 接入层路由 │ ├── core/ # 配置、日志、异常 │ ├── models/ # Pydantic 模型和 ORM 模型 │ ├── services/ # 聚合层业务逻辑 │ ├── collectors/ # 采集层适配器 │ └── storage/ # 存储层封装 ├── config/ │ ├── mappings/ # 数据源字段映射 │ └── settings.yaml ├── data/ # Parquet 文件存储 ├── tests/ └── scripts/ # 运维脚本3.2 核心 API 的实现与参数设计对外暴露的 API 我设计得很克制只保留最必要的几个端点。端点太多会增加维护负担而且大部分需求可以通过参数组合满足。核心端点就四个GET /api/v1/bars查询 K 线数据GET /api/v1/symbols查询标的列表GET /api/v1/calendar查询交易日历GET /api/v1/health健康检查/api/v1/bars的参数设计我反复调整过好几版最终定下来这几个参数类型必填说明symbolstring是标的代码如 000001.SZfreqstring是频率1m/5m/15m/30m/60m/1d/1w/1Mstartdate是开始日期ISO 格式enddate是结束日期ISO 格式adjuststring否复权方式none/qfq/hfq默认 qfqfieldsstring否逗号分隔的字段列表默认全部limitint否最大返回条数默认 5000上限 50000为什么把limit上限设成 50000因为再大的话单次响应体可能超过 10MB网络传输和客户端解析都会变慢。需要更多数据就分页拉或者直接用批量导出接口。响应格式统一成 JSON结构如下{ code: 0, message: ok, data: { symbol: 000001.SZ, freq: 1d, adjust: qfq, bars: [ {ts: 2024-01-02T00:00:00Z, open: 10.5, high: 10.8, low: 10.3, close: 10.6, volume: 1234567} ] }, meta: { count: 1, source: aggregated, cached: false } }code为 0 表示成功非 0 表示各类错误。错误码我定义了一套简单的规则1xxx 是参数错误2xxx 是数据不存在3xxx 是上游异常5xxx 是服务内部错误。3.3 聚合层的对齐与计算逻辑聚合层是整套服务的灵魂。它要做的事情是从多个源拿到数据后决定用哪个源的数据、怎么对齐、怎么计算衍生指标。源选择策略我实现了一个简单的优先级加健康度评分机制。每个源有一个基础优先级分数再根据最近的成功率、延迟、数据完整度动态调整。每次查询时选综合得分最高的源。如果最高分源返回的数据不完整自动降级到次高分源补全。def select_source(sources: list, symbol: str, freq: str) - str: scored [] for src in sources: health get_health_score(src.name, symbol, freq) score src.priority * 0.6 health * 0.4 scored.append((src.name, score)) scored.sort(keylambda x: x[1], reverseTrue) return scored[0][0]复权计算是另一个核心逻辑。前复权qfq和后复权hfq的公式不同但原理都是基于除权除息日的调整因子。我维护一张corporate_actions表记录每次分红送配的细节然后根据查询区间计算累积调整因子。前复权的调整因子计算方式是从查询区间末尾往前累乘每次除权除息事件对应一个因子。后复权则是从区间开头往后累乘。具体公式这里不展开网上资料很多关键是理解“前复权看近期价格真实后复权看历史收益真实”这个本质。注意复权计算一定要用精确的除权除息数据不能用近似值。我见过有人用“分红率”粗略估算结果长周期回测的收益率偏差能到 5% 以上策略直接失效。3.4 定时任务与数据同步调度数据同步我用 APScheduler 做调度没有引入 Celery 这种重家伙。原因很简单我的同步任务数量不多并发要求也不高APScheduler 足够用而且部署简单不用额外维护消息队列。调度策略分三档日频数据每个交易日收盘后 30 分钟触发拉取当天数据分钟频数据交易时段内每 5 分钟触发一次增量拉取元数据标的列表、交易日历每天凌晨 2 点全量刷新from apscheduler.schedulers.background import BackgroundScheduler from apscheduler.triggers.cron import CronTrigger scheduler BackgroundScheduler(timezoneUTC) scheduler.add_job( sync_daily_bars, CronTrigger(day_of_weekmon-fri, hour8, minute30), idsync_daily, max_instances1, misfire_grace_time3600 ) scheduler.add_job( sync_intraday_bars, CronTrigger(day_of_weekmon-fri, hour1-7,9-15, minute*/5), idsync_intraday, max_instances1 )max_instances1这个参数很关键防止上一次任务还没跑完下一次就启动了导致数据重复写入。misfire_grace_time设 3600 秒意思是如果任务因为服务重启错过了触发时间一小时内还会补跑。4. 常见问题与排查技巧实录4.1 数据源相关的典型故障做金融数据服务跟上游数据源打交道的时间占了一大半。下面这几个问题我几乎每个源都遇到过。问题一接口返回 200 但数据为空。这种情况最坑因为 HTTP 层面看不出任何异常。我的排查步骤是先看请求参数是否被上游静默忽略再看返回体里有没有error字段但被放在了非标准位置最后对比同一时间其他源的数据确认是源的问题还是标的的问题。解决方式是在 adapter 里加一层“空数据告警”连续三次返回空就自动标记该源为不健康。问题二字段类型突变。有的源平时返回数字某天突然返回字符串比如10.5而不是10.5。Pydantic 默认会尝试强制转换但有时候转不了就报错。我的做法是在 normalize 阶段统一做类型转换用pd.to_numeric(errorscoerce)把转不了的变成 NaN然后走缺失值处理流程。问题三频率漂移。分钟频数据偶尔会多出几秒或者少几秒导致时间戳对不齐。解决方式是在聚合层做时间桶对齐把所有时间戳向下取整到最近的 bar 边界。问题现象可能原因排查方法解决方案返回空数据源限流/标的停牌/参数错误对比多源查日志加空数据告警自动降级字段类型突变上游改接口抽样检查类型校验normalize层强制转换时间戳漂移上游时钟不准统计时间戳分布时间桶对齐数据重复重试机制bug查sync_state表幂等写入唯一索引复权因子错误除权数据缺失对比前后复权价格补全corporate_actions4.2 性能瓶颈的定位与优化服务上线初期我遇到过一次查询超时。单个请求要 8 秒以上用户体验极差。排查过程记录如下。第一步加日志打点。在 API 入口、聚合层入口、存储层入口分别记录时间戳算出各阶段耗时。结果发现 90% 的时间花在存储层查询上。第二步看慢查询。PostgreSQL 的pg_stat_statements显示有个查询在扫全表。原因是sync_state表没建索引每次查询都在做 seq scan。加上(source_name, symbol, freq)的联合索引后查询时间从 6 秒降到 50 毫秒。第三步看缓存命中率。Redis 的INFO stats显示命中率只有 40%远低于预期的 80%。原因是缓存 key 里包含了精确到秒的时间戳导致同样的查询每次 key 都不同。改成按天粒度生成 key 之后命中率升到 85%。第四步看 Parquet 文件。发现有些文件特别大单个超过 1GB读取时内存直接爆掉。改成按月分区后单文件控制在 200MB 以内配合 DuckDB 的流式读取内存占用稳定在 500MB 以下。实操心得性能优化一定要先测量再动手。我一开始凭感觉优化了半天代码逻辑结果发现瓶颈根本不在那里。加日志打点虽然土但真的管用。4.3 数据一致性保障的独家技巧金融数据最怕的就是不一致。同一份数据今天查和明天查结果不一样那所有基于它的分析都不可信。我用了几个手段来保障一致性。第一个手段是写入幂等。所有数据写入操作都基于(symbol, freq, ts)这个唯一键做 upsert重复写入不会产生重复记录。Parquet 文件写入时先写临时文件写完校验通过再原子重命名避免写到一半崩溃导致文件损坏。第二个手段是版本快照。每天同步完成后对当天的数据生成一个校验和存到data_versions表里。如果后续发现数据有问题可以追溯到具体是哪天的同步出了错。import hashlib def compute_checksum(df) - str: content df.sort_values(ts).to_csv(indexFalse) return hashlib.sha256(content.encode()).hexdigest()第三个手段是对账机制。每周跑一次全量对账把本地数据和上游数据做抽样比对发现偏差超过阈值就告警。抽样比例我设的是 5%既能发现问题又不至于太耗资源。4.4 安全与权限的底线配置虽然是个人项目但安全底线不能破。我做了几件事。API 鉴权用简单的 API Key 机制每个 Key 绑定一个调用配额。Key 存在 PostgreSQL 里用 bcrypt 哈希后存储不存明文。请求头里带X-API-Key中间件校验通过才放行。限流用 Redis 的滑动窗口算法每个 Key 每分钟最多 60 次请求。超过就返回 429并在响应头里带上Retry-After。def check_rate_limit(api_key: str, limit: int 60, window: int 60) - bool: now int(time.time()) key fratelimit:{api_key}:{now // window} pipe redis.pipeline() pipe.incr(key) pipe.expire(key, window * 2) count, _ pipe.execute() return count limit数据库连接串、API Key 这些敏感信息全部走环境变量不写进代码库。.env文件加到.gitignore里防止误提交。注意即使是内部服务也不要用0.0.0.0监听所有网卡。我一开始图方便这么干后来发现同一网络下的其他设备能直接访问赶紧改成只监听127.0.0.1外部访问走反向代理。5. 后续扩展方向与个人体会这套服务跑了大半年目前支撑着我自己的几个策略回测和两个朋友的记账工具日调用量在几千次左右单次查询 P99 延迟控制在 200 毫秒以内。过程中最大的体会是金融数据服务的难点不在技术而在对数据本身的理解。同样的 K 线数据复权方式不同、时间对齐方式不同、缺失值处理方式不同算出来的结果可能天差地别。技术只是工具真正决定服务质量的是你对业务规则的理解深度。后续我打算加两个东西。一个是 WebSocket 推送把实时行情推给订阅的客户端省掉轮询开销。另一个是数据质量看板把每个源的可用性、延迟、完整度可视化出来方便快速定位问题。这两个都不难等有空了慢慢做。如果你也在搭类似的东西我的建议是先把采集和存储跑通哪怕只接一个源、只存一种频率。等数据能稳定落盘了再往上叠聚合和 API。不要一上来就追求大而全那样大概率会烂尾。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门