当你的策略开始赚钱,问题才刚刚开始
凌晨两点,你从床上爬起来查看策略表现。
不是因为时差,不是因为焦虑。而是你的笔记本跑着全部策略,本地 SQLite 数据库已经膨胀到 8GB,每次回测都要等上十分钟。更要命的是,室友上个月不小心把你写的新因子覆盖了——他只是想“同步一下代码”。
这不是段子。这是每一个从小做大做大的量化团队都会经历的阵痛。
一个人做量化,瓶颈是策略本身。三个人的量化团队,瓶颈是基础设施——数据怎么共享、代码怎么管理、密钥怎么流转、权限怎么控制。这些问题不会让你的策略亏钱,但会偷走你所有的时间。
本文拆解一个 3-5 人小型量化团队如何从“每人一套孤岛”走向生产级协作架构。
一、小团队的数据基础设施有哪些典型死法
在谈方案之前,先承认现实。多数小团队的崩溃不是一夜之间发生的,而是三件事反复叠加的结果。
死法一:每人一套数据,自己回测自己
策略研究员 A 用他的历史数据回测,收益曲线漂亮得离谱。实盘三个月,收益腰斩。问题不在策略,在于他的数据和实盘数据差了 12 个百分点——他用的是未复权的行情,而实盘系统接的是前复权数据。
数据源不统一,回测结果就是废纸。
死法二:Git 成了噩梦
策略研究员 B 提交代码从来不写 commit message,"wip"写了 80 次。策略研究员 C 发现自己分支落后了 200 次提交,直接删掉本地重拉——连带着把 B 的未合并因子一起删了。
没有分支策略,没有 code review,没有强制 commit 规范,Git 就是一颗定时炸弹。
死法三:API Key 裸奔在代码里
三年前写的脚本里,api_key = "sk-xxxxx" 直接写死在代码里。三年后实习生拿到了代码仓库权限,连带着把账户额度烧光了。
没有密钥管理机制,所有的安全措施都是纸糊的。
死法四:权限控制等于没有控制
junior 获得了生产数据库的写权限,一次 UPDATE 语句没有加 WHERE,半年的因子数据全变成 NULL。
权限的松紧程度决定了事故的破坏半径。
这四件事不会同时爆发,但会逐个累积。当团队终于开始赚钱、需要扩容的时候,你会发现重建基础设施的成本比继续缝缝补补还要高。
二、生产级架构:四层结构
小型量化团队的基础设施不需要 Kubernetes,不需要微服务架构,但需要四层清晰的分离:
┌─────────────────────────────────────────────────────┐
│ 接入层 │
│ (TickDB / 数据源统一封装) │
├─────────────────────────────────────────────────────┤
│ 存储层 │
│ (时序数据库 + 文件存储 + 配置存储) │
├─────────────────────────────────────────────────────┤
│ 协作层 │
│ (Git + CI/CD + 环境管理) │
├─────────────────────────────────────────────────────┤
│ 安全层 │
│ (密钥管理 + 权限控制 + 操作审计) │
└─────────────────────────────────────────────────────┘
每层各司其职,层与层之间通过接口而非实现交互。这是团队基础设施能持续演进的核心。
三、接入层:数据源统一封装
3.1 为什么需要统一封装
如果团队里每个人直接调用原始 API,数据会迅速变得不可控:
- 某人的脚本硬编码了
symbol = "BTCUSD",另一个人写symbol = "BTC-USD" - 有人在循环里同步调用,有人用异步
- 当数据源 API 变更时,每处调用都要改
统一封装的本质是定义团队的 API 标准。
3.2 封装设计原则
# src/data/adapters/base.py
from abc import ABC, abstractmethod
from typing import List, Optional, Dict, Any
from dataclasses import dataclass
from datetime import datetime
import os
@dataclass
class MarketData:
"""统一市场数据结构"""
symbol: str
timestamp: datetime
open: float
high: float
low: float
close: float
volume: float
source: str # 标记数据来源,便于溯源
@dataclass
class OrderBook:
"""统一订单簿结构"""
symbol: str
timestamp: datetime
bids: List[tuple] # [(price, volume), ...]
asks: List[tuple]
source: str
class BaseDataAdapter(ABC):
"""
数据适配器基类
所有数据源必须实现此接口,确保:
1. 返回类型统一
2. 错误处理一致
3. 日志可追溯
"""
def __init__(self, api_key: str):
self.api_key = api_key
self._rate_limit_state = {"last_call": 0, "call_count": 0}
@abstractmethod
def get_kline(
self,
symbol: str,
interval: str,
start_time: Optional[datetime] = None,
end_time: Optional[datetime] = None,
limit: int = 1000
) -> List[MarketData]:
"""获取 K 线数据"""
pass
@abstractmethod
def get_depth(self, symbol: str, limit: int = 10) -> OrderBook:
"""获取订单簿深度"""
pass
def _check_rate_limit(self, interval_seconds: int = 1):
"""简单限频检查"""
import time
current_time = time.time()
if current_time - self._rate_limit_state["last_call"] < interval_seconds:
self._rate_limit_state["call_count"] += 1
if self._rate_limit_state["call_count"] > 10:
raise RuntimeError("Rate limit exceeded")
else:
self._rate_limit_state["call_count"] = 0
self._rate_limit_state["last_call"] = current_time
3.3 TickDB 适配器实现
# src/data/adapters/tickdb_adapter.py
import os
import time
import random
import json
import logging
from typing import List, Optional, Dict, Any
from datetime import datetime
from .base import BaseDataAdapter, MarketData, OrderBook
logger = logging.getLogger(__name__)
class TickDBAdapter(BaseDataAdapter):
"""
TickDB 数据源适配器
⚠️ 生产环境高频场景建议使用 aiohttp/asyncio
特性:
- 支持美股、港股、数字货币、外汇等多品种
- K 线数据 10 年级别历史覆盖
- depth 频道最大支持 50 档深度
"""
BASE_URL = "https://api.tickdb.ai/v1"
def __init__(self, api_key: Optional[str] = None):
# 从环境变量读取,永远不硬编码
actual_key = api_key or os.environ.get("TICKDB_API_KEY")
if not actual_key:
raise ValueError("TickDB API Key 未配置。请设置环境变量 TICKDB_API_KEY")
super().__init__(actual_key)
def _request_with_retry(
self,
method: str,
endpoint: str,
params: Optional[Dict] = None,
max_retries: int = 3
) -> Dict:
"""
带重试的 HTTP 请求
实现要点:
- 指数退避 + 抖动,避免惊群效应
- 识别限频错误码 3001,等待 Retry-After
- 超时设置防止挂起
"""
import requests
url = f"{self.BASE_URL}{endpoint}"
headers = {"X-API-Key": self.api_key}
base_delay = 1
max_delay = 32
for attempt in range(max_retries):
try:
response = requests.request(
method=method,
url=url,
headers=headers,
params=params,
timeout=(3.05, 10) # (connect_timeout, read_timeout)
)
# 识别限频错误
if response.status_code == 429 or (
response.headers.get("Content-Type", "").startswith("application/json")
and response.json().get("code") == 3001
):
retry_after = int(response.headers.get("Retry-After", 5))
logger.warning(f"限频触发,等待 {retry_after} 秒")
time.sleep(retry_after)
continue
response.raise_for_status()
result = response.json()
if result.get("code") != 0:
raise RuntimeError(f"API Error: {result.get('message')}")
return result.get("data", {})
except requests.exceptions.Timeout:
logger.warning(f"请求超时,重试 ({attempt + 1}/{max_retries})")
except requests.exceptions.RequestException as e:
logger.warning(f"请求异常: {e},重试 ({attempt + 1}/{max_retries})")
# 指数退避 + 抖动
delay = min(base_delay * (2 ** attempt), max_delay)
jitter = random.uniform(0, delay * 0.1)
time.sleep(delay + jitter)
raise RuntimeError(f"请求失败,已重试 {max_retries} 次")
def get_kline(
self,
symbol: str,
interval: str,
start_time: Optional[datetime] = None,
end_time: Optional[datetime] = None,
limit: int = 1000
) -> List[MarketData]:
"""获取 K 线数据"""
params = {
"symbol": symbol,
"interval": interval,
"limit": limit
}
if start_time:
params["start_time"] = int(start_time.timestamp())
if end_time:
params["end_time"] = int(end_time.timestamp())
self._check_rate_limit()
raw_data = self._request_with_retry("GET", "/market/kline", params)
results = []
for item in raw_data.get("klines", []):
results.append(MarketData(
symbol=symbol,
timestamp=datetime.fromtimestamp(item["timestamp"]),
open=float(item["open"]),
high=float(item["high"]),
low=float(item["low"]),
close=float(item["close"]),
volume=float(item["volume"]),
source="TickDB"
))
return results
def get_depth(self, symbol: str, limit: int = 10) -> OrderBook:
"""获取订单簿深度(实时接口示例见 WebSocket 部分)"""
params = {"symbol": symbol, "limit": limit}
self._check_rate_limit()
raw_data = self._request_with_retry("GET", "/market/depth", params)
return OrderBook(
symbol=symbol,
timestamp=datetime.now(),
bids=[(float(p), float(v)) for p, v in raw_data.get("bids", [])],
asks=[(float(p), float(v)) for p, v in raw_data.get("asks", [])],
source="TickDB"
)
# src/data/registry.py
class DataAdapterRegistry:
"""
数据适配器注册表
团队所有数据源在此注册,确保:
- 同一数据源全局单例
- 方便切换数据源进行对比
"""
_adapters: Dict[str, BaseDataAdapter] = {}
@classmethod
def register(cls, name: str, adapter: BaseDataAdapter):
cls._adapters[name] = adapter
logger.info(f"注册数据适配器: {name}")
@classmethod
def get(cls, name: str) -> BaseDataAdapter:
if name not in cls._adapters:
raise KeyError(f"未找到数据适配器: {name}")
return cls._adapters[name]
@classmethod
def initialize_from_config(cls, config_path: str = "config/adapters.yaml"):
"""从配置文件初始化所有适配器"""
import yaml
if not os.path.exists(config_path):
logger.warning(f"配置文件不存在: {config_path}")
return
with open(config_path) as f:
config = yaml.safe_load(f)
for adapter_name, adapter_config in config.get("adapters", {}).items():
if adapter_config["type"] == "tickdb":
adapter = TickDBAdapter()
cls.register(adapter_name, adapter)
3.4 团队使用规范
# 团队统一的数据获取方式
from src.data.registry import DataAdapterRegistry
def get_close_prices(symbols: List[str], interval: str = "1d", limit: int = 100):
"""团队统一的数据获取入口"""
adapter = DataAdapterRegistry.get("tickdb")
results = {}
for symbol in symbols:
try:
klines = adapter.get_kline(symbol, interval, limit=limit)
results[symbol] = [k.close for k in klines]
except Exception as e:
logger.error(f"获取 {symbol} 数据失败: {e}")
results[symbol] = None
return results
这样无论谁写回测脚本,数据来源完全一致,回测结果可以对比。
四、存储层:数据分层管理
4.1 存储分层设计
小型团队的存储不需要分布式文件系统,但需要三层分离:
| 层级 | 存储类型 | 工具选择 | 用途 |
|---|---|---|---|
| 时序数据 | 时序数据库 | TimescaleDB / InfluxDB | K 线、因子、订单簿快照 |
| 文件存储 | 对象存储 | S3 / MinIO | 代码、模型、配置文件 |
| 配置存储 | 版本化配置 | Git + YAML | 环境变量、策略参数 |
4.2 时序数据存储设计
-- 创建时序表(TimescaleDB 示例)
CREATE TABLE market_kline (
time TIMESTAMPTZ NOT NULL,
symbol TEXT NOT NULL,
interval TEXT NOT NULL,
open NUMERIC(18, 8),
high NUMERIC(18, 8),
low NUMERIC(18, 8),
close NUMERIC(18, 8),
volume NUMERIC(18, 8),
source TEXT -- 数据来源标记:AAPL 是从 TickDB 获取
);
SELECT create_hypertable('market_kline', 'time');
-- 创建索引加速查询
CREATE INDEX idx_kline_symbol_interval ON market_kline (symbol, interval, time DESC);
-- 保留策略:美股数据保留 10 年,数字货币保留 5 年
SELECT add_retention_policy('market_kline', INTERVAL '10 years')
WHERE symbol ~ '^(AAPL|MSFT|NVDA)';
4.3 统一存储服务封装
# src/storage/unified_store.py
from contextlib import contextmanager
from typing import Generator
import logging
from sqlalchemy import create_engine
from sqlalchemy.orm import sessionmaker, Session
from sqlalchemy.pool import QueuePool
logger = logging.getLogger(__name__)
class UnifiedStorage:
"""
统一存储服务
⚠️ 每次操作后必须关闭 session,避免连接泄漏
"""
_instance = None
def __new__(cls):
if cls._instance is None:
cls._instance = super().__new__(cls)
cls._instance._initialized = False
return cls._instance
def initialize(self, database_url: str):
if self._initialized:
return
# 连接池配置:小型团队 5-10 个连接足够
self.engine = create_engine(
database_url,
poolclass=QueuePool,
pool_size=5,
max_overflow=10,
pool_pre_ping=True, # 连接前检测
pool_recycle=3600 # 1 小时回收
)
self.SessionLocal = sessionmaker(
autocommit=False,
autoflush=False,
bind=self.engine
)
self._initialized = True
logger.info("存储服务初始化完成")
@contextmanager
def session(self) -> Generator[Session, None, None]:
"""安全的 session 管理"""
session = self.SessionLocal()
try:
yield session
session.commit()
except Exception as e:
session.rollback()
logger.error(f"数据库操作失败: {e}")
raise
finally:
session.close()
# 全局存储实例
storage = UnifiedStorage()
def init_storage():
db_url = os.environ.get("DATABASE_URL")
if not db_url:
raise ValueError("DATABASE_URL 未配置")
storage.initialize(db_url)
五、协作层:Git 规范化
5.1 分支策略
小型团队的 Git 流程不需要 GitFlow 那么重,但需要三条硬规则:
main ────────────────────────────────────→ (生产环境)
↑ merge
│
release/X.Y ←────────────────┘ (发布准备)
↑ merge from feature
│
feature/xxx ──→ (开发中)
↑
┌────┴────┐
dev-a dev-b (个人分支,daily sync)
强制规则:
main分支禁止直接 push,必须通过 PR- 每次 PR 必须关联一个 Issue
- Commit message 必须符合 Conventional Commits 规范
5.2 Commit Message 规范
<type>(<scope>): <subject>
<body>
<footer>
| Type | 使用场景 |
|---|---|
| feat | 新功能 |
| fix | Bug 修复 |
| docs | 文档更新 |
| refactor | 重构(不改变功能) |
| test | 测试相关 |
| chore | 工具、依赖更新 |
好的 commit message:
feat(kline): 添加前复权数据转换支持
- 新增 adjust_kline() 函数处理复权因子
- 回测模块默认使用前复权数据
- 修复因子计算中浮点精度问题
Closes #42
坏的 commit message:
wip
fix
update
5.3 Git Hooks 强制执行
# .git/hooks/commit-msg
#!/bin/bash
commit_msg=$(cat "$1")
# 检查是否符合 Conventional Commits
if ! echo "$commit_msg" | grep -qE '^(feat|fix|docs|refactor|test|chore)(\(.+\))?: .+'; then
echo "❌ Commit message 格式错误"
echo "正确格式: type(scope): subject"
echo "例如: feat(kline): 添加新因子"
exit 1
fi
echo "✅ Commit message 格式检查通过"
# .git/hooks/pre-push
#!/bin/bash
branch=$(git symbolic-ref --short HEAD)
# 禁止直接 push 到 main
if [ "$branch" = "main" ]; then
echo "❌ 禁止直接 push 到 main 分支"
echo "请创建 PR 并请求 code review"
exit 1
fi
echo "✅ 分支检查通过"
5.4 代码审查清单
每个 PR 必须满足以下条件才能合并:
| 检查项 | 说明 |
|---|---|
| 测试覆盖 | 新功能必须有单元测试 |
| Lint 通过 | ruff check 无警告 |
| 文档更新 | 改动影响 API 必须更新文档 |
| 无敏感信息 | 确认不包含 API Key 或密码 |
六、安全层:密钥与权限
6.1 密钥管理原则
核心原则:密钥永远不进代码仓库。
❌ 错误做法
api_key = "sk-xxxxx" # 直接写死在代码里
✅ 正确做法
api_key = os.environ.get("TICKDB_API_KEY") # 从环境变量读取
6.2 密钥管理架构
# .env.example # 提交到仓库,仅作为模板
# 团队成员复制此文件为 .env.local,不提交到仓库
TICKDB_API_KEY=your_api_key_here
DATABASE_URL=postgresql://user:pass@localhost:5432/quant
S3_ACCESS_KEY=your_access_key
S3_SECRET_KEY=your_secret_key
# .gitignore # 确保敏感文件被忽略
.env.local
.env.*.local
*.pem
*.key
secrets/
6.3 环境配置加载器
# src/config/loader.py
import os
import logging
from pathlib import Path
from typing import Any, Dict
from dotenv import load_dotenv
logger = logging.getLogger(__name__)
def load_environment():
"""
加载环境配置
优先级(从高到低):
1. 系统环境变量
2. .env.local
3. .env
"""
# 查找 .env 文件
env_files = [
Path.cwd() / ".env.local",
Path(__file__).parent.parent.parent / ".env.local",
Path.cwd() / ".env",
]
for env_file in env_files:
if env_file.exists():
load_dotenv(env_file, override=False)
logger.info(f"加载环境配置: {env_file}")
break
# 检查必需的环境变量
required_vars = ["TICKDB_API_KEY"]
missing = [v for v in required_vars if not os.environ.get(v)]
if missing:
logger.warning(f"未配置的环境变量: {', '.join(missing)}")
logger.warning("请参考 .env.example 配置后重试")
class Config:
"""配置单例"""
_instance = None
_config: Dict[str, Any] = {}
def __new__(cls):
if cls._instance is None:
cls._instance = super().__new__(cls)
load_environment()
return cls._instance
@property
def tickdb_api_key(self) -> str:
key = os.environ.get("TICKDB_API_KEY")
if not key:
raise EnvironmentError("TICKDB_API_KEY 未配置")
return key
@property
def database_url(self) -> str:
return os.environ.get(
"DATABASE_URL",
"postgresql://localhost:5432/quant"
)
6.4 数据库权限控制
-- 创建角色(按职责分离)
CREATE ROLE quant_reader; -- 只读
CREATE ROLE quant_writer; -- 可读写因子和信号
CREATE ROLE quant_admin; -- 管理权限
-- 授予权限
GRANT SELECT ON ALL TABLES IN SCHEMA public TO quant_reader;
GRANT SELECT, INSERT, UPDATE ON market_kline TO quant_writer;
GRANT SELECT, INSERT, UPDATE, DELETE ON factors TO quant_writer;
GRANT ALL ON ALL TABLES IN SCHEMA public TO quant_admin;
-- junior 成员只能使用 quant_reader,避免误操作
-- 策略研究员使用 quant_writer
-- 团队负责人使用 quant_admin
七、团队部署方案对比
| 方案 | 适用规模 | 部署成本 | 维护复杂度 | 推荐指数 |
|---|---|---|---|---|
| 全手动 | 1-2 人 | 低 | 高 | ⭐ |
| Docker Compose | 3-5 人 | 中 | 中 | ⭐⭐⭐⭐ |
| Kubernetes | 5 人以上 | 高 | 高 | ⭐⭐ |
| 云托管服务 | 任意规模 | 中-高 | 低 | ⭐⭐⭐ |
推荐 3-5 人团队使用 Docker Compose,平衡了部署便利性和资源占用。
# docker-compose.yml
version: '3.8'
services:
postgres:
image: timescale/timescaledb:latest-pg15
environment:
POSTGRES_DB: quant
POSTGRES_USER_FILE: /run/secrets/db_user
POSTGRES_PASSWORD_FILE: /run/secrets/db_pass
volumes:
- postgres_data:/var/lib/postgresql/data
secrets:
- db_user
- db_pass
minio:
image: minio/minio
command: server /data --console-address ":9001"
environment:
MINIO_ROOT_USER_FILE: /run/secrets/s3_user
MINIO_ROOT_PASSWORD_FILE: /run/secrets/s3_pass
volumes:
- minio_data:/data
secrets:
- s3_user
- s3_pass
volumes:
postgres_data:
minio_data:
secrets:
db_user:
file: ./secrets/db_user.txt
db_pass:
file: ./secrets/db_pass.txt
s3_user:
file: ./secrets/s3_user.txt
s3_pass:
file: ./secrets/s3_pass.txt
八、从零到一的检查清单
如果你是 3 人团队准备搭建基础设施,按这个顺序执行:
第一周:统一数据源
- 选定数据源(TickDB 推荐)
- 实现统一适配器
- 配置 CI 自动同步历史数据到本地数据库
- 团队成员统一切换到新适配器
第二周:代码规范
- 搭建 Git 仓库,配置 hooks
- 制定 commit message 规范
- 建立 code review 流程
- 完成第一轮 code review
第三周:密钥与环境
- 搭建密钥管理机制
- 完成 .env 配置
- 配置数据库权限
- 文档化部署流程
第四周:运维自动化
- Docker Compose 编排
- 自动备份策略
- 监控告警配置
- 灾难恢复演练
结语
基础设施不是“等团队大了再搞”的事情。
小团队的基础设施问题,等于策略的机会成本。每次回测等 10 分钟,每年就是 60 个小时。每次数据不一致导致的策略失效,都是真金白银的损失。
搭建基础设施最好的时机是现在——团队还小,改造成本还低。
第二好的时机是爆雷之后。
下一步行动
如果你是 1-2 人团队,当前数据获取还是“各自为战”:
- 访问 tickdb.ai 了解统一数据接入方案
- 下载本文的代码模板,在团队内部署适配器
- 用一周时间完成数据源统一
如果你已经有基本架构,但缺乏规范:
- 从 Git hooks 强制 commit message 格式开始
- 用 Docker Compose 重构部署脚本
- 配置数据库分角色权限
如果你希望团队数据基础设施从零开始:
联系 [email protected],获取针对量化团队的专项技术方案支持。
风险提示:本文不构成任何投资建议。基础设施搭建需要结合团队实际情况评估,文中代码仅供参考,生产环境部署请进行充分测试。