1. 金融数据服务从零搭建的核心思路拆解
1.1 为什么我要自己动手做一套金融数据服务
先说清楚这套东西是干什么的。financial-services,直译就是“金融服务”,但在我这里,它指的是一套面向个人开发者和小型团队的自建金融数据服务层。它能做什么?简单讲,就是把行情数据、财报数据、宏观经济指标这些散落在各处的金融信息,通过一套统一的接口聚合起来,对外提供稳定的查询、计算和推送能力。解决的核心问题是:数据源太散、格式不统一、调用方式五花八门,每次做个策略回测或者搭个看板,光数据对接就要耗掉大半精力。
这套内容适合谁?如果你是一个做量化策略的个人开发者,或者是一个小团队里负责数据基建的工程师,再或者你是一个对金融数据感兴趣、想自己搭一套工具链的技术爱好者,那这篇内容就是写给你的。我不假设你有海量服务器资源,也不假设你有昂贵的商业数据终端,所有的方案都是基于常见公开数据源和常规技术栈来设计的,追求的是“够用、稳定、可维护”。
为什么我要强调“自建”而不是“直接用第三方API”?这里有个很现实的考量。市面上的金融数据接口,免费的限制多、延迟高,付费的又贵得离谱,而且一旦你的业务逻辑深度绑定某一家,后面想换源就是伤筋动骨。自建服务层的价值在于,它在你和原始数据源之间加了一个抽象层。上层业务只跟你的统一接口打交道,底层换数据源、加缓存、做降级,对上层都是透明的。这个设计思路,跟微服务里“网关”的角色很像——把复杂性收敛到一个地方,让其他地方保持简单。
1.2 整体架构选型:为什么是“采集-清洗-存储-服务”四层
我把整个financial-services拆成四层,这个分层不是拍脑袋定的,而是踩过坑之后总结出来的。最早我试过“直连模式”,业务代码里直接调数据源API,结果就是每个用到数据的地方都要写一遍重试、解析、异常处理,代码重复率极高,而且数据源一改字段,满项目找地方改。后来改成两层,采集和服务混在一起,又发现采集的调度逻辑和服务的查询逻辑互相干扰,一个慢查询能把采集任务拖死。
最终定下来的四层结构是这样的:
- 采集层:负责从各个数据源拉取原始数据,只做最基础的格式校验,不做复杂转换。这一层的核心任务是“把数据拿回来”,不关心数据怎么用。
- 清洗层:对原始数据做标准化处理,包括字段重命名、单位统一、时间对齐、缺失值处理。这一层是保证数据质量的关键。
- 存储层:根据数据的访问模式选择不同的存储介质。时序数据用时序库,关系型数据用关系库,需要全文检索的用搜索引擎。
- 服务层:对外提供统一的RESTful接口和WebSocket推送,处理鉴权、限流、缓存、降级等逻辑。
这个分层的优势在于,每一层都可以独立扩展和替换。比如采集层今天用A数据源,明天想换成B,只要输出格式不变,上层完全无感。清洗层的规则可以随时调整,不影响采集的稳定性。存储层可以根据数据量增长平滑迁移。服务层则是对外的门面,所有的SLA保障都在这一层做。
注意:分层不是目的,解耦才是。如果你的数据量很小、业务逻辑很简单,硬套四层反而是过度设计。我建议先从两层做起,等痛了再拆。
1.3 技术栈选择背后的逻辑
技术栈这块,我选型的原则是“成熟优先、社区活跃、运维成本低”。具体来说:
- 采集层:Python + APScheduler。Python的生态在金融数据处理这块确实丰富,pandas、numpy这些库处理数值计算很顺手。APScheduler做定时任务调度,轻量够用,支持cron表达式和间隔触发。为什么不用Airflow?Airflow太重了,对于个人项目来说,运维一个Airflow实例的成本可能比业务本身还高。
- 清洗层:还是Python,用pandas做DataFrame级别的转换。pandas的向量化操作在处理批量数据时性能很好,而且语法直观,写清洗规则很快。
- 存储层:PostgreSQL + TimescaleDB扩展 + Redis。PostgreSQL存关系型数据和元数据,TimescaleDB是PostgreSQL的时序扩展,存行情数据非常合适,支持自动分区和压缩。Redis做缓存和消息队列,用来缓冲采集任务和加速热点查询。
- 服务层:FastAPI + Uvicorn。FastAPI的异步性能好,自动生成OpenAPI文档,类型提示友好,开发效率高。Uvicorn作为ASGI服务器,部署简单。
这套组合的总体思路是:用最少的组件覆盖最多的场景,避免引入不必要的复杂度。每一个组件都是经过大量项目验证的,遇到问题容易找到解决方案。
2. 核心细节解析与实操要点
2.1 数据源接入的标准化封装
数据源接入是采集层的核心工作。我接触过的金融数据源大概分三类:HTTP API、文件下载、数据库直连。不管哪一类,我都会把它封装成一个统一的“DataSource”抽象类,定义几个必须实现的方法:
from abc import ABC, abstractmethod from typing import List, Dict, Any class DataSource(ABC): @abstractmethod def fetch(self, symbol: str, start: str, end: str) -> List[Dict[str, Any]]: """拉取指定标的在时间范围内的数据""" pass @abstractmethod def health_check(self) -> bool: """检查数据源是否可用""" pass @property @abstractmethod def source_name(self) -> str: """数据源名称,用于日志和监控""" pass这样封装的好处是,每个数据源的差异被限制在各自的实现类里。比如某个数据源返回的是JSON,另一个返回的是CSV,但在fetch方法内部都转换成统一的List[Dict]格式。上层清洗层拿到的数据格式是一致的,不需要关心数据从哪来。
实操中有一个细节很重要:每个数据源都要有独立的限流和重试策略。不同数据源的QPS限制不一样,有的严格有的宽松。我会在DataSource的实现类里内置一个令牌桶限流器,根据数据源的实际承受能力配置速率。重试策略用指数退避,但要注意区分“可重试错误”和“不可重试错误”。比如网络超时可以重试,但参数错误重试多少次都没用。
提示:数据源的健康检查不要只检查网络连通性,要实际拉一条数据验证返回格式是否符合预期。我遇到过数据源接口没挂但返回结构变了的情况,只检查连通性会漏掉这种问题。
2.2 数据清洗的规则设计与实现
清洗层的核心任务是“把原始数据变成可信数据”。我总结了几条必须处理的规则:
字段标准化。不同数据源对同一个字段的叫法可能完全不同。比如开盘价,有的叫open,有的叫open_price,有的叫开盘价。我会维护一个字段映射表,把所有变体映射到统一的标准字段名。这个映射表用YAML配置,方便修改。
时间对齐。金融数据对时间非常敏感。不同数据源的时间戳可能是UTC,可能是北京时间,可能带时区信息,可能不带。我的做法是统一转换成UTC时间戳存储,在服务层再根据用户请求的时区做转换。另外,对于日线数据,要明确时间戳代表的是“交易日开始”还是“交易日结束”,这个不统一会导致回测结果完全错误。
缺失值处理。金融数据缺失很常见,可能是数据源本身的问题,也可能是非交易日。我的策略是:先区分“真缺失”和“假缺失”。非交易日的数据缺失是正常的,不应该填充。真缺失则根据字段类型处理:价格类字段用前值填充,成交量类字段用0填充,但都要打上标记,方便后续分析时过滤。
异常值检测。金融数据里偶尔会出现离谱的异常值,比如价格突然变成0或者负数。我会用简单的统计方法做检测:计算滚动窗口的均值和标准差,超出3倍标准差的标记为异常。异常值不直接删除,而是标记出来,由人工确认后再决定处理方式。
import pandas as pd import numpy as np def clean_ohlcv(df: pd.DataFrame) -> pd.DataFrame: df = df.copy() # 字段标准化 df = df.rename(columns=FIELD_MAPPING) # 时间对齐 df['timestamp'] = pd.to_datetime(df['timestamp'], utc=True) # 异常值标记 for col in ['open', 'high', 'low', 'close']: rolling_mean = df[col].rolling(window=20, min_periods=5).mean() rolling_std = df[col].rolling(window=20, min_periods=5).std() df[f'{col}_anomaly'] = np.abs(df[col] - rolling_mean) > 3 * rolling_std # 缺失值处理 df['close'] = df['close'].ffill() df['volume'] = df['volume'].fillna(0) return df2.3 存储层的表结构设计与索引优化
存储层这块,我踩过的最大坑是“一张表存所有”。最早我把所有标的的行情数据放在一张表里,数据量上来之后查询慢得没法用。后来改成按标的和时间分区,性能提升非常明显。
具体来说,行情数据表的设计是这样的:
CREATE TABLE ohlcv ( symbol VARCHAR(20) NOT NULL, timestamp TIMESTAMPTZ NOT NULL, open NUMERIC(18, 6), high NUMERIC(18, 6), low NUMERIC(18, 6), close NUMERIC(18, 6), volume NUMERIC(24, 6), PRIMARY KEY (symbol, timestamp) ); SELECT create_hypertable('ohlcv', 'timestamp', chunk_time_interval => INTERVAL '1 month');TimescaleDB的create_hypertable会自动按时间分区,每个月的分区独立存储和索引。查询时如果带上时间范围条件,TimescaleDB会自动裁剪掉不相关的分区,扫描的数据量大幅减少。
索引方面,主键(symbol, timestamp)已经覆盖了“查某个标的某段时间”这个最常见的查询模式。如果还需要按时间范围查所有标的,可以再加一个timestamp的索引。但索引不是越多越好,每个索引都会增加写入开销。我的原则是:只为实际用到的查询模式建索引,不确定的先不建,等慢了再加。
Redis的用法主要是两个场景:一是缓存热点查询结果,比如最新行情、常用标的的基本信息;二是作为采集任务的消息队列,采集层把任务推到Redis,清洗层从Redis消费。用Redis做队列的好处是轻量,不需要额外部署消息中间件。但要注意Redis的持久化配置,如果对任务可靠性要求高,要开启AOF。
2.4 服务层接口设计与限流降级
服务层是对外的门面,设计好坏直接影响使用体验。我的接口设计遵循几个原则:
统一响应格式。所有接口返回统一的JSON结构,包含code、message、data三个字段。code为0表示成功,非0表示各种错误。这样客户端处理起来逻辑统一。
分页与游标。对于可能返回大量数据的接口,必须支持分页。我倾向于用游标分页而不是偏移量分页,因为偏移量分页在数据量大时性能差,而且数据变动时会出现重复或遗漏。游标分页用时间戳或自增ID作为游标,性能稳定。
限流策略。服务层必须做限流,否则一个异常客户端就能把整个服务打挂。我用的是基于Redis的滑动窗口限流,按API Key维度限制每分钟的请求数。限流阈值根据实际承载能力设定,留出足够的余量。
降级方案。当后端数据源不可用时,服务层不能直接报错,而是要有降级策略。我的做法是:优先返回缓存数据,缓存也没有则返回最近一次成功获取的数据并标记stale字段,同时触发告警。这样至少保证客户端能拿到数据,而不是完全不可用。
from fastapi import FastAPI, HTTPException, Depends from fastapi.responses import JSONResponse app = FastAPI() @app.get("/api/v1/ohlcv") async def get_ohlcv(symbol: str, start: str, end: str, api_key: str = Depends(verify_api_key)): # 限流检查 if not rate_limiter.allow(api_key): raise HTTPException(status_code=429, detail="Rate limit exceeded") # 尝试从缓存获取 cached = cache.get(f"ohlcv:{symbol}:{start}:{end}") if cached: return {"code": 0, "message": "ok", "data": cached} # 从数据库查询 try: data = query_ohlcv(symbol, start, end) cache.set(f"ohlcv:{symbol}:{start}:{end}", data, ttl=60) return {"code": 0, "message": "ok", "data": data} except Exception as e: # 降级:返回最近一次成功的数据 fallback = cache.get(f"ohlcv:fallback:{symbol}") if fallback: return {"code": 0, "message": "degraded", "data": fallback, "stale": True} raise HTTPException(status_code=503, detail="Service unavailable")3. 实操过程与核心环节实现
3.1 环境搭建与依赖安装
环境搭建这块,我推荐用Docker Compose来管理所有依赖服务,这样环境一致性好,迁移也方便。下面是我用的docker-compose.yml核心部分:
version: '3.8' services: postgres: image: timescale/timescaledb:latest-pg15 environment: POSTGRES_DB: financial POSTGRES_USER: fin_user POSTGRES_PASSWORD: fin_pass ports: - "5432:5432" volumes: - pg_data:/var/lib/postgresql/data redis: image: redis:7-alpine ports: - "6379:6379" command: redis-server --appendonly yes volumes: - redis_data:/data volumes: pg_data: redis_data:Python依赖用requirements.txt管理:
fastapi==0.104.1 uvicorn[standard]==0.24.0 pandas==2.1.3 numpy==1.26.2 psycopg2-binary==2.9.9 redis==5.0.1 apscheduler==3.10.4 httpx==0.25.1 pydantic==2.5.2 pyyaml==6.0.1安装命令很简单:
pip install -r requirements.txt这里有个细节要注意:psycopg2-binary和psycopg2的区别。binary版本是预编译的,安装快,适合开发和测试。生产环境建议用源码编译的psycopg2,性能和稳定性更好。另外,如果用的是Apple Silicon的Mac,某些包的wheel可能不兼容,需要从源码编译,这时候要确保Xcode Command Line Tools已安装。
3.2 采集任务的调度与执行
采集任务的调度我用APScheduler,配置如下:
from apscheduler.schedulers.blocking import BlockingScheduler from apscheduler.triggers.cron import CronTrigger scheduler = BlockingScheduler() # 日线数据:每个交易日收盘后采集 scheduler.add_job( collect_daily_ohlcv, CronTrigger(day_of_week='mon-fri', hour=18, minute=0), id='daily_ohlcv', max_instances=1, misfire_grace_time=3600 ) # 实时行情:交易时段每5秒采集一次 scheduler.add_job( collect_realtime_quote, CronTrigger(day_of_week='mon-fri', hour='9-15', minute='*', second='*/5'), id='realtime_quote', max_instances=1 ) scheduler.start()max_instances=1这个参数很重要,它保证同一个任务不会并发执行。如果上一次采集还没跑完,下一次触发会被跳过,避免数据重复或资源竞争。misfire_grace_time是错过触发的宽限时间,比如服务器重启导致任务错过了,在宽限时间内还会补跑一次。
采集任务的执行逻辑我封装成一个装饰器,统一处理日志、重试和异常:
import functools import time import logging logger = logging.getLogger(__name__) def with_retry(max_retries=3, backoff_base=2): def decorator(func): @functools.wraps(func) def wrapper(*args, **kwargs): for attempt in range(max_retries): try: return func(*args, **kwargs) except RetryableError as e: if attempt == max_retries - 1: logger.error(f"{func.__name__} failed after {max_retries} attempts: {e}") raise wait = backoff_base ** attempt logger.warning(f"{func.__name__} attempt {attempt+1} failed, retrying in {wait}s") time.sleep(wait) except NonRetryableError as e: logger.error(f"{func.__name__} failed with non-retryable error: {e}") raise return wrapper return decorator这个装饰器的关键点是区分了RetryableError和NonRetryableError。网络超时、服务暂时不可用属于可重试;参数错误、数据格式不匹配属于不可重试。这个区分很重要,否则会在无意义的错误上浪费大量重试时间。
3.3 数据清洗流水线的搭建
清洗流水线我用的是“管道-过滤器”模式,每个清洗步骤是一个独立的函数,数据依次流过这些函数。这样做的好处是每个步骤可以单独测试,也可以灵活组合。
from typing import Callable, List import pandas as pd class CleaningPipeline: def __init__(self): self.steps: List[Callable[[pd.DataFrame], pd.DataFrame]] = [] def add_step(self, step: Callable[[pd.DataFrame], pd.DataFrame]) -> 'CleaningPipeline': self.steps.append(step) return self def run(self, df: pd.DataFrame) -> pd.DataFrame: for step in self.steps: df = step(df) return df # 使用示例 pipeline = (CleaningPipeline() .add_step(standardize_fields) .add_step(align_timestamps) .add_step(detect_anomalies) .add_step(handle_missing) .add_step(validate_schema)) cleaned_df = pipeline.run(raw_df)每个清洗步骤的编写有几个要点。第一,步骤函数必须是纯函数,不修改输入DataFrame,而是返回新的DataFrame。这样便于调试和回滚。第二,每个步骤都要有日志记录,记录处理前后的行数变化和关键字段的统计信息。第三,步骤的顺序很重要,比如字段标准化必须在时间对齐之前,因为时间对齐可能依赖标准化的字段名。
validate_schema这一步我强烈建议加上。它检查清洗后的数据是否符合预期的schema,包括字段是否存在、类型是否正确、取值范围是否合理。这一步能在数据进入存储层之前拦截大部分问题,避免脏数据污染数据库。
3.4 服务接口的部署与监控
服务层用Uvicorn部署,生产环境建议用Gunicorn做进程管理,配合Uvicorn Worker:
gunicorn main:app \ --workers 4 \ --worker-class uvicorn.workers.UvicornWorker \ --bind 0.0.0.0:8000 \ --timeout 120 \ --access-logfile - \ --error-logfile -workers的数量一般是CPU核心数的2倍加1。但金融数据服务通常是IO密集型的,等待数据库和外部API的时间占大头,所以可以适当增加worker数量。timeout设置要合理,太短会导致长查询被中断,太长会导致僵尸进程堆积。
监控这块,我用Prometheus + Grafana。FastAPI有现成的prometheus-fastapi-instrumentator库,几行代码就能接入:
from prometheus_fastapi_instrumentator import Instrumentator Instrumentator().instrument(app).expose(app)关键监控指标包括:请求量、响应时间P50/P95/P99、错误率、缓存命中率、数据库连接池使用率。这些指标能帮你快速定位性能瓶颈。比如P99响应时间突然飙升,可能是某个查询没走索引;缓存命中率下降,可能是缓存key设计有问题或者TTL设置太短。
提示:监控告警的阈值不要设得太敏感,否则会被大量误报淹没。我一般先观察一周的正常波动范围,再根据P99的1.5倍设置告警阈值。
4. 常见问题与排查技巧实录
4.1 数据源不稳定导致采集失败
这是最常见的问题。公开数据源经常出现超时、限流、返回格式变化等情况。我的排查思路是分三步走:
第一步,确认是网络问题还是数据源问题。用curl或httpx直接请求数据源,看返回状态码和响应时间。如果直接请求也失败,那就是数据源的问题;如果直接请求成功但程序里失败,那可能是程序配置或代码问题。
第二步,检查限流和重试配置。很多数据源对请求频率有严格限制,超过就返回429。这时候要检查令牌桶的速率设置是否合理,重试策略是否生效。我遇到过重试间隔太短导致连续触发限流的情况,把退避基数从1秒改成2秒就解决了。
第三步,验证返回数据格式。数据源可能在不通知的情况下调整返回结构,比如字段改名、嵌套层级变化。这时候清洗层的validate_schema会报错,根据错误信息定位到具体字段,更新字段映射表即可。
| 问题现象 | 可能原因 | 排查方法 | 解决方案 |
|---|---|---|---|
| 请求超时 | 网络抖动或数据源过载 | 直接curl测试 | 增加超时时间,加重试 |
| 返回429 | 触发限流 | 检查请求频率 | 降低采集频率,加令牌桶 |
| 字段缺失 | 数据源格式变化 | 对比返回结构 | 更新字段映射表 |
| 数据重复 | 任务并发执行 | 检查调度配置 | 设置max_instances=1 |
| 数据延迟 | 数据源更新慢 | 对比其他数据源 | 调整采集时间或换源 |
4.2 数据库查询性能突然下降
数据库查询变慢通常有几个原因:数据量增长导致索引失效、查询计划变化、连接池耗尽、锁竞争。我的排查顺序是:
先看慢查询日志。PostgreSQL的pg_stat_statements扩展能记录所有SQL的执行统计,按平均执行时间排序,很快就能找到最慢的查询。找到慢查询后,用EXPLAIN ANALYZE看执行计划,重点关注是否有全表扫描、是否走了错误的索引。
如果是数据量增长导致的,考虑加索引或者调整分区策略。TimescaleDB的分区是按时间自动创建的,但如果单个分区数据量还是太大,可以缩小chunk_time_interval,比如从1个月改成1周。
连接池耗尽也是常见问题。FastAPI默认没有连接池管理,每次请求都新建数据库连接,并发一高就撑不住。解决方案是用asyncpg的连接池或者SQLAlchemy的QueuePool。连接池大小根据实际并发量设置,一般max_size设为worker数量的2到3倍。
from sqlalchemy import create_engine from sqlalchemy.pool import QueuePool engine = create_engine( "postgresql://fin_user:fin_pass@localhost:5432/financial", poolclass=QueuePool, pool_size=10, max_overflow=20, pool_pre_ping=True, pool_recycle=3600 )pool_pre_ping=True会在每次从池中取连接时先ping一下,避免使用已失效的连接。pool_recycle=3600让连接每小时回收一次,防止数据库端主动断开。
4.3 缓存与数据库数据不一致
缓存不一致是分布式系统的经典问题。我的策略是“缓存失效优先于缓存更新”。也就是说,当数据更新时,不是去更新缓存,而是直接删除缓存,下次查询时自然从数据库加载最新数据并重建缓存。
这个策略的逻辑是:更新缓存和更新数据库是两个操作,无论谁先谁后,都有不一致的窗口。而删除缓存只有一个操作,不一致的窗口更小。具体实现:
def update_ohlcv(symbol, timestamp, data): # 先更新数据库 db.execute("UPDATE ohlcv SET ... WHERE symbol=%s AND timestamp=%s", ...) # 再删除缓存 cache.delete(f"ohlcv:{symbol}:{timestamp}") # 删除相关列表缓存 cache.delete_pattern(f"ohlcv:list:{symbol}:*")delete_pattern要慎用,如果key数量太多会阻塞Redis。更好的做法是用版本号或者命名空间来管理缓存,更新时只递增版本号,旧版本的缓存自然失效。
注意:如果对一致性要求极高,可以考虑用“延迟双删”——更新数据库后删一次缓存,延迟几百毫秒再删一次。这是为了应对“更新数据库后、删除缓存前”有读请求把旧数据写回缓存的情况。但大多数金融数据场景对实时一致性要求没那么高,简单删除就够了。
4.4 服务接口被恶意刷量
服务层暴露在公网,难免会遇到恶意刷量。除了前面提到的限流,还有几个防护措施:
API Key鉴权。每个客户端分配一个API Key,请求时带上。API Key可以设置权限范围和配额。没有Key或者Key无效的请求直接拒绝。
IP黑名单。对于频繁触发限流的IP,加入黑名单,一段时间内拒绝所有请求。黑名单用Redis的Set实现,设置过期时间自动解除。
请求签名。对于敏感接口,要求客户端对请求参数做签名,服务端验证签名。这样能防止参数被篡改,也能增加刷量的成本。
监控告警。当某个API Key的请求量突增,或者错误率突增时,触发告警。我一般设置两个阈值:请求量超过日常均值3倍,或者错误率超过10%。
def verify_api_key(api_key: str = Header(...)): key_info = cache.get(f"apikey:{api_key}") if not key_info: raise HTTPException(status_code=401, detail="Invalid API key") if key_info['quota_exceeded']: raise HTTPException(status_code=429, detail="Quota exceeded") return key_info这套防护措施下来,一般的恶意刷量都能挡住。但如果遇到大规模的DDoS,那就不是应用层能解决的了,需要在上层做流量清洗。
4.5 数据回补与历史数据修复
数据采集难免会有遗漏,比如某天服务器宕机导致数据没采到,或者发现历史数据有错误需要修复。这时候就需要数据回补机制。
我的做法是维护一个“数据完整性检查”任务,每天跑一次,检查最近N天的数据是否有缺失。发现缺失就生成回补任务,推入队列。回补任务和正常采集任务走同一套清洗和存储流程,保证数据一致性。
def check_data_integrity(symbol: str, start: str, end: str): expected_dates = get_trading_dates(start, end) actual_dates = query_existing_dates(symbol, start, end) missing_dates = set(expected_dates) - set(actual_dates) for date in missing_dates: queue.push({'symbol': symbol, 'date': date, 'type': 'backfill'})回补任务要注意限流,不能因为回补把正常采集的配额用完了。我会给回补任务设置更低的优先级和更保守的速率限制。另外,回补的数据要标记来源,方便后续审计。
历史数据修复更麻烦一些,因为涉及已存储数据的更新。我的原则是:不直接修改原始数据,而是把修正后的数据作为新版本写入,查询时取最新版本。这样保留了数据变更历史,出问题可以追溯。
5. 个人实操心得与后续扩展方向
5.1 几个让我少走弯路的经验
第一,日志要打够,但不要打太多。我早期为了排查问题,把每个请求的完整参数和返回都打到日志里,结果日志文件一天几十个G,磁盘经常满。后来改成:正常请求只记录关键字段和耗时,异常请求记录完整上下文。日志级别用INFO,DEBUG只在排查特定问题时临时开启。
第二,配置和代码分离。所有可能变化的参数,比如数据源地址、限流阈值、缓存TTL,都放到配置文件或环境变量里。这样改配置不需要改代码,也不需要重新部署。我用YAML做配置文件,用Pydantic做配置校验,启动时如果配置有问题直接报错,避免运行到一半才发现。
第三,先跑通再优化。我见过太多人(包括我自己)一开始就追求完美的架构,结果花了大量时间在设计上,真正跑起来发现根本不是那么回事。我的建议是:先用最简单的方式跑通整个流程,哪怕采集是手动触发的、存储是CSV文件、服务是Flask单进程。跑通之后,根据实际遇到的瓶颈逐步优化。这样每一步优化都有明确的收益,不会过度设计。
第四,数据质量比数据量重要。金融数据里,一条错误的数据可能比没有数据更糟糕。我宁愿少采一些数据,也要保证采到的数据是准确的。所以清洗层的validate_schema和异常值检测是必须的,不能省。
5.2 这套服务还能怎么扩展
当前这套financial-services覆盖了行情数据的基本场景,但金融数据的范围远不止这些。后续可以扩展的方向包括:
财报数据接入。财报数据的结构和行情数据完全不同,需要单独设计表结构和清洗规则。财报数据的特点是低频、结构化程度高、字段多。可以考虑用JSONB字段存储原始财报,用关系型字段存储关键指标。
宏观经济指标。这类数据通常来自统计部门或国际组织,更新频率低但历史数据长。存储上可以用单独的表,按指标类型分区。
实时推送服务。当前的服务层主要是请求-响应模式,如果要支持实时推送,需要引入WebSocket。FastAPI原生支持WebSocket,但要注意连接管理和消息广播的效率。连接数多的时候,单机WebSocket可能撑不住,需要考虑用Redis Pub/Sub做消息分发。
回测引擎集成。数据服务的最终目的往往是支撑策略回测。可以在服务层之上再加一个回测引擎,直接调用数据服务获取历史数据,运行策略逻辑,输出回测报告。这样整个链路就完整了。
多租户支持。如果这套服务要给多个用户使用,需要加租户隔离。数据层面可以用schema隔离或者行级权限,接口层面用API Key关联租户ID,所有查询自动带上租户过滤条件。
这套东西我从最初的一个脚本,慢慢迭代到现在这个规模,前后大概花了半年时间。中间踩过的坑、推翻重来的设计不在少数。但每次优化都是被实际问题驱动的,所以每一步都走得踏实。如果你也在做类似的事情,我的建议是:不要追求一步到位,先让数据流跑起来,然后在实践中不断打磨。数据服务这个东西,稳定性和准确性永远是第一位的,花哨的功能反而是其次。