当你的策略开始赚钱,问题才刚刚开始

凌晨两点,你从床上爬起来查看策略表现。

不是因为时差,不是因为焦虑。而是你的笔记本跑着全部策略,本地 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 人团队,当前数据获取还是“各自为战”

  1. 访问 tickdb.ai 了解统一数据接入方案
  2. 下载本文的代码模板,在团队内部署适配器
  3. 用一周时间完成数据源统一

如果你已经有基本架构,但缺乏规范

  1. 从 Git hooks 强制 commit message 格式开始
  2. 用 Docker Compose 重构部署脚本
  3. 配置数据库分角色权限

如果你希望团队数据基础设施从零开始
联系 [email protected],获取针对量化团队的专项技术方案支持。


风险提示:本文不构成任何投资建议。基础设施搭建需要结合团队实际情况评估,文中代码仅供参考,生产环境部署请进行充分测试。