给中级以上工程师 / 架构评审 / 跨团队协作的一份深度文档。 聚焦"每个子包内部是什么样、和主项目怎么协作"——在飞书云文档里直接粘贴可用。
- 子包全景
- quantcore — 纯量化计算核
- tickbridge — 统一数据层
- querybus — 统一搜索层
- pulsefan — 统一通知层
- stocklens — 业务编排层
- api / bot — 外部适配层
- 子包间协作时序
- 一图看懂数据走向
StockLens 整个项目由 5 个独立可发布的核心包 + 2 个适配层 构成:
| 层级 | 包名 | 职责 | 行数 | 文件数 | 独立发布 |
|---|---|---|---|---|---|
| 算法层 | quantcore/ |
纯计算(评分/指标/统计/回测) | ~4,400 | 19 | ✅ |
| 数据层 | tickbridge/ |
行情 / 基本面 / 筹码统一抓取 | ~17,200 | 64 | ✅ |
| 搜索层 | querybus/ |
新闻 / 搜索结果统一抓取 | ~3,100 | 32 | ✅ |
| 通知层 | pulsefan/ |
多渠道并行推送 | ~3,500 | 30 | ✅ |
| 业务层 | stocklens/ |
编排上面 4 层 + 业务逻辑 | ~40,500 | 175 | ❌ |
| 适配层 | api/ |
FastAPI REST | ~6,300 | 28 | ❌ |
| 适配层 | bot/ |
钉钉/飞书 Stream Bot | ~2,800 | 16 | ❌ |
| 入口 | main.py |
启动 FastAPI + 可选 Bot | 159 | 1 | — |
┌─────────┬─────────┬─────────┬─────────┐
│quantcore│tickbridge│querybus │pulsefan │ ← 4 个独立包,互不引用
└────┬────┴────┬────┴────┬────┴────┬────┘
↑ ↑ ↑ ↑
└─────────┴─────────┴─────────┘
│
┌────┴────┐
│stocklens│ ← 业务编排层
└────┬────┘
│
┌────┴────┐
│ │
api/ bot/ ← 适配层
│ │
└────┬────┘
│
main.py
铁律:
- quantcore / tickbridge / querybus / pulsefan 互相之间不能 import
- 它们 不能 import stocklens / api / bot
- stocklens 不能 import api / bot
- 跨层共享类型 → 必须放在
stocklens.contracts/
输入 DataFrame / dict,输出 dataclass。零 IO,零网络,零数据库。
- ❌ 不准 import
requests/akshare/sqlite3/ 任何文件 IO - ❌ 不准 import 其他业务包
- ✅ 仅依赖
numpy+pandas - ✅ 需要外部数据时,通过参数注入(caller 抓好了传进来)
quantcore/
├── __init__.py 暴露 ScoringEngine + BacktestEngine 等顶层类
├── pyproject.toml 独立包定义
│
├── scoring/ ★ 18 量化评分模型
│ ├── types.py 18 个 dataclass + classify_industry / _sf / _unwrap_block
│ ├── fundamental.py Piotroski / Altman Z / Ohlson O / Composite Distress
│ ├── valuation.py Rule of 40 / PEG / FCF Yield / Graham / Beneish M /
│ │ DuPont / EarningsQuality / Magic Formula / DividendSafety
│ ├── sentiment.py Insider / Options / Momentum / EarningsSurprise / Short Interest
│ ├── technical.py SCTR / Risk / Macro / MarketRegime / PeerRank
│ └── engine.py ScoringEngine — 编排 18 模型 + composite_score(0-100)
│
├── indicators/ ★ 纯技术指标
│ ├── trend.py StockTrendAnalyzer / TrendSnapshot
│ │ (MA / MACD / RSI / BB Squeeze / RS vs SPY)
│ ├── bollinger.py BollingerResult / MultiBollingerSnapshot + compute_bollinger
│ ├── resample.py 通用 OHLCV resample 原语(W/ME/QE/YE/Nh)
│ ├── labels.py describe_volume_ratio / compute_ma_status
│ └── tech.py FibonacciLevels / KlinePattern / TechIndicatorSnapshot
│ + detect_swing_points / select_fibonacci_swing
│ + compute_tech_indicators(Fibonacci/ATR/OBV/VWAP/K线/Sharpe/MaxDD)
│
├── stats/ ★ 统计/风险原语
│ ├── risk.py compute_sharpe / compute_sortino / compute_max_drawdown / compute_volatility
│ └── correlation.py compute_correlation_pair(成对 stock/benchmark → CorrelationResult)
│
└── backtest/ ★ 回测评估引擎
└── engine.py BacktestEngine(中英双语关键词 + 否定识别 + 止损止盈命中 + 胜率聚合)
EvaluationConfig + DailyBarLike / BacktestResultLike Protocol
# 评分
from quantcore.scoring import (
ScoringEngine, ScoringOverview,
calc_piotroski, calc_distress,
calc_valuation, calc_beneish, calc_dupont,
calc_options_sentiment, calc_momentum,
calc_sctr, calc_risk, calc_macro, calc_market_regime,
classify_industry,
)
# 指标
from quantcore.indicators import (
StockTrendAnalyzer, TrendSnapshot,
BollingerResult, MultiBollingerSnapshot, compute_bollinger,
TechIndicatorSnapshot, FibonacciLevels, compute_tech_indicators,
resample_ohlcv, resample_to_weekly, resample_to_monthly,
resample_to_quarterly, resample_to_yearly, resample_hourly,
describe_volume_ratio, compute_ma_status,
)
# 统计
from quantcore.stats import (
compute_sharpe, compute_sortino, compute_max_drawdown, compute_volatility,
compute_correlation_pair,
)
# 回测
from quantcore.backtest import (
BacktestEngine, EvaluationConfig,
DailyBarLike, BacktestResultLike, OVERALL_SENTINEL_CODE,
)stocklens 把 IO 包在外面,把纯计算丢给 quantcore:
| stocklens 文件 | 留下的 IO 部分 | 委托给 quantcore 的部分 |
|---|---|---|
stocklens.services.scoring.ScoringService |
注入 CBOE 期权抓取 fetcher(网络调用) | quantcore.scoring.ScoringEngine |
stocklens.services.bollinger_service.fetch_bollinger_multi_timeframe |
akshare/雪球抓取 + 二级缓存 | quantcore.indicators.compute_bollinger |
stocklens.services.tech_indicators_service.compute_correlation |
akshare 抓 SPY/QQQ + 4h 缓存 | quantcore.stats.compute_correlation_pair |
stocklens.backtest.engine |
纯 re-export shim(旧导入路径仍可用) | quantcore.backtest.engine |
stocklens.pipeline.data_collector 静态 helper |
业务编排 | quantcore.indicators.{resample, labels} |
所有评分都用 Python 算,LLM 不可修改。
为什么这么设计?避免 AI 幻觉污染数据。LLM 只在最后一步"解读已经算好的分数+指标",分数本身永远可复现。
# 实战调用流程(stocklens 视角)
from quantcore.scoring import ScoringEngine
from stocklens.services.scoring._cboe_fetcher import fetch_cboe_options
engine = ScoringEngine(cboe_fetcher=fetch_cboe_options)
overview = engine.evaluate(context) # 跑 18 个模型
score = ScoringEngine.calc_composite_score(overview) # 0-100 综合分
table = ScoringEngine.format_overview_table(overview) # 给报告用的 markdown 表格所有"获取行情/基本面/筹码"的请求都走这个包,10 个数据源用统一 API 接入,自动 failover。
- ❌
api / bot / stocklens都不准直接from tickbridge.sources.X import Y - ❌
api / bot / stocklens都不准requests.get任何上游 API(雪球/Finviz/Reddit 等) - ✅ 唯一入口:
provider = DataProvider(); provider.get_realtime_quote(...) - ✅ 后端特定方法:
provider.source('xueqiu').fetch_bars(...)
tickbridge/
├── __init__.py 暴露 DataProvider + __version__
├── pyproject.toml 独立包定义
├── LICENSE / README.md MIT
│
├── core/ ★ 门面 / 调度器 / 三级读路径
│ ├── provider.py DataProvider(唯一对外类)
│ ├── manager.py DataFetcherManager — 多源 failover 调度
│ ├── circuit_breaker.py per-source 熔断器(失败冷却)
│ ├── read_through.py cache → DB → source 三级读路径
│ ├── realtime_mixin.py 实时报价能力
│ ├── chip_mixin.py 筹码分布
│ ├── fundamental_mixin.py 基本面聚合
│ ├── market_mixin.py 板块 / 指数 / 市场统计
│ ├── stock_info_mixin.py 股票名 / sector
│ ├── exceptions.py DataFetchError / RateLimitError
│ └── utils.py safe_float / safe_int
│
├── contracts/ BaseFetcher / DomainFetcher / CacheProtocol / RepositoryProtocol
├── models/ 公开数据契约
│ ├── quote.py UnifiedRealtimeQuote
│ ├── chip.py ChipDistribution
│ └── ohlcv / fundamental / market
│
├── sources/ ★ 10 个 BaseFetcher 子类(多源 failover)
│ ├── yfinance.py P0 Yahoo Finance
│ ├── xueqiu.py P1 雪球
│ ├── tushare.py P2 Tushare
│ ├── openbb.py P2 OpenBB
│ ├── efinance.py P4 东方财富
│ ├── akshare.py P5 AKShare
│ ├── pytdx.py P7 通达信协议
│ ├── baostock.py P8 Baostock
│ ├── twelvedata.py P999 Twelve Data
│ └── tickflow.py TickFlow
│
├── domain_sources/ ★ 单源 DomainFetcher(领域专用工具)
│ ├── finviz.py FinvizFetcher(前瞻数据 / 同行 / 英文新闻 / 内部交易)
│ └── social_sentiment.py SocialSentimentService(Reddit / X 情绪)
│
├── codes/ 股票代码工具
│ ├── normalize.py normalize_stock_code / 市场识别
│ ├── us_index.py US 指数符号映射
│ └── stock_mapping.py 代码 → 中文名
│
├── adapters/ 上游 → 公开模型字段适配
├── internal/ 私有运行时补丁(如东方财富 NID 授权)
├── config/ DI 钩子(register_config_provider)
├── cache/ 默认缓存实现
└── storage/ 默认 SQLite Repository
provider = DataProvider()
# 实时行情
quote = provider.get_realtime_quote("AAPL")
quote = provider.get_index_quote(".INX")
quote = provider.get_etf_quote("SPY")
# 历史 K 线
df = provider.get_daily_data("AAPL", days=200)
# 筹码分布(A 股)
chip = provider.get_chip_distribution("000001")
# 基本面
fund = provider.get_fundamental("AAPL")
# 板块 / 指数 / 市场
indices = provider.get_main_indices(region="us")
stats = provider.get_market_stats()
sectors = provider.get_sector_rankings(n=5)
# Finviz 领域工具(仅美股)
fwd = provider.finviz_forward("AAPL")
news = provider.finviz_news("AAPL", limit=15)
peers = provider.finviz_peers("AAPL", max_peers=8)
insider = provider.finviz_insider("AAPL")
# 社交情绪(仅美股)
reddit = provider.social_reddit_report("AAPL")
x_trending = provider.social_x_trending()get_realtime_quote("AAPL")
│
▼
1️⃣ in-memory TTLCache(命中?返回)
│ MISS
▼
2️⃣ SQLite cache(命中且未过期?返回)
│ MISS / EXPIRED
▼
3️⃣ DataFetcherManager.fetch():
│
├─ Yahoo (P0, 优先级最高) ──→ 失败(熔断器记录 1 次失败)
│
├─ Xueqiu (P1) ──→ 成功 ✓
│
├─ Tushare (P2) ──→ 跳过(已有结果)
│ ...
│
└─ Twelve Data (P999, 兜底)
│
▼
写回 SQLite cache + in-memory
│
▼
返回结果
每个数据源独立计数失败:
- 连续失败 3 次 → 熔断 60 秒(不再尝试)
- 熔断到期 → 半开(试探一次)
- 半开成功 → 完全恢复
- 半开失败 → 重新熔断 120 秒(指数退避)
# stocklens 业务代码(绝对不会出现 import 雪球的代码)
from tickbridge import DataProvider
class DataCollector:
def __init__(self, fetcher_manager):
self.provider = DataProvider() # 由 DI 注入也可
def fetch_and_save(self, code):
# 一行调用,背后跑完三级读路径
quote = self.provider.get_realtime_quote(code)
df = self.provider.get_daily_data(code, days=200)
# 写入 stocklens 自己的 SQLite ORM
self.stock_repo.save(quote, df)所有"搜新闻 / 搜公告 / 搜社交舆情"的请求都走这个包,7 个搜索引擎统一接入。
- ❌
api / bot / stocklens不准直接from querybus.providers.X import Y - ❌
api / bot / stocklens不准requests.get调 Tavily / Brave / Bocha 等任何搜索 API - ✅ 唯一入口:
provider = SearchProvider(...); provider.search(...) - ✅ 股票领域包装层在
stocklens.search.StockSearchService
querybus/
├── __init__.py 暴露 SearchProvider + SearchResult + SearchResponse + DimensionSpec
├── pyproject.toml 独立包定义
│
├── core/ ★ 门面 / 调度器 / 缓存 / 过滤
│ ├── provider.py SearchProvider(唯一对外类)
│ ├── manager.py SearchManager — 多源 fall-over
│ ├── cache.py InMemoryTTLCache(FIFO + TTL)
│ ├── filters.py filter_by_time_window + dedupe_results
│ └── reporting.py format_intel_report + DimensionSpec
│
├── contracts/ BaseSearchProvider / CacheProtocol / retry
├── models/
│ ├── result.py SearchResult(单条结果)
│ └── response.py SearchResponse(一次查询的完整响应)
│
├── providers/ ★ 7 个 BaseSearchProvider 子类
│ ├── tavily.py Tavily AI 搜索
│ ├── serpapi.py SerpAPI(Google 代理)
│ ├── bocha.py 博查(中文)
│ ├── brave.py Brave Search
│ ├── minimax.py MiniMax 搜索
│ ├── searxng.py SearXNG(自部署)
│ └── akshare_news.py AKShare 财经新闻
│
└── parsing/ 解析工具
├── relative_dates.py "3 天前" / "yesterday"
├── absolute_dates.py "2026-05-20"
└── content.py lazy newspaper3k 正文抽取
from querybus import SearchProvider
provider = SearchProvider()
# 普通搜索
resp = provider.search("Apple Q4 earnings", max_results=20, days_back=7)
# 批量搜索
responses = provider.batch_search(
queries=["Apple stock", "AAPL fundamentals"],
max_results=10,
)
# 后端特定调用(绕过 fall-over)
provider.provider("tavily").search("...")
# 后处理工具
filtered = SearchProvider.filter_by_time_window(resp.results, days=3)
SearchProvider.dedupe(responses)
# 渲染情报报告
report = SearchProvider.format_intel_report(
responses=multi_dim_responses,
dimensions=[DimensionSpec(name="财报", ...), ...]
)search(query="Apple Q4 earnings", days_back=7)
│
▼
1️⃣ 内存 TTL 缓存(hash(query, days_back) 作为 key)
│ MISS
▼
2️⃣ 按优先级遍历 7 个 provider:
Tavily → SerpAPI → Bocha → Brave → MiniMax → SearXNG → AKShare News
│
├─ 成功 → 时间窗过滤 → 去重 → 写缓存 → 返回
└─ 全部失败 → 抛 SearchExhaustedError
querybus 本身只懂"通用搜索"。"搜某只股票的新闻"这个领域语义放在 stocklens 侧的薄壳:
# stocklens/search/stock_search_service.py
from querybus import SearchProvider
class StockSearchService:
def __init__(self):
self.provider = SearchProvider(...)
def search_stock_news(self, code, name):
# 拼股票领域 query:name + 新闻关键词 + 时间窗
query = f"{name} 股票 财报 业绩 OR earnings"
return self.provider.search(query, days_back=14)
def gather_intel(self, code, name):
# 5 维度并发抓取(财报/同行/政策/技术/资金)
with ThreadPoolExecutor(max_workers=5) as ex:
...所有"推送报告 / 发通知"的请求都走这个包,10 个渠道并行扇出,一个失败不影响其他。
- ❌
api / bot / stocklens不准直接from pulsefan.senders.X import Y - ❌
api / bot / stocklens不准requests.post调微信/飞书/Telegram 任何推送 API - ❌ 不准
from stocklens.utils.formatters import X(已迁移到pulsefan.formatting) - ✅ 唯一入口:
notifier = Notifier(config); notifier.send_text(...) - ✅ 业务报告编排在
stocklens.notification.ReportOrchestrator
pulsefan/
├── __init__.py 暴露 Notifier + NotificationConfig + 10 ChannelConfig + SendResult
├── pyproject.toml
│
├── core/ ★ 门面 / 注册表 / 并行扇出
│ ├── notifier.py Notifier(唯一对外类)
│ ├── registry.py ChannelRegistry — 从 NotificationConfig 实例化 senders
│ ├── dispatcher.py ParallelDispatcher — ThreadPoolExecutor fan-out
│ └── result.py SendResult dataclass
│
├── contracts/
│ └── base_sender.py BaseSender (ABC)
│
├── config/
│ └── schema.py NotificationConfig + 10 个 ChannelConfig dataclass
│ (pulsefan 自己不读环境变量,调用方传 dataclass 进来)
│
├── senders/ ★ 10 个 BaseSender 子类
│ ├── wechat.py 企业微信
│ ├── feishu.py 飞书
│ ├── telegram.py Telegram
│ ├── email.py 邮件 SMTP
│ ├── discord.py Discord
│ ├── pushover.py Pushover
│ ├── pushplus.py PushPlus
│ ├── serverchan3.py Server 酱 v3
│ ├── astrbot.py AstrBot
│ └── custom_webhook.py 自定义 webhook
│
└── formatting/ ★ 纯 Python 工具(不发请求)
├── chunking.py 长文本分块(超长消息切片)
├── markdown.py markdown → 纯文本 / HTML
├── feishu.py 飞书富文本/卡片格式
└── _legacy_formatters.py 旧版兼容
from pulsefan import Notifier, NotificationConfig, FeishuConfig, EmailConfig
cfg = NotificationConfig(
feishu=FeishuConfig(webhook="https://open.feishu.cn/...", secret="..."),
email=EmailConfig(host="smtp.qq.com", user="...", password="..."),
)
notifier = Notifier(cfg)
# 文本推送(所有已配置渠道并行)
results = notifier.send_text(title="盘后报告", content="...")
# results: List[SendResult],逐渠道成功/失败
# 图片推送
notifier.send_image(image_path="report.png", caption="日报")
# 单渠道精确控制(绕过并行扇出)
notifier.sender("wechat").send_to_wechat(title="...", content="...")
# 查询配置状态
notifier.configured_channels() # ['feishu', 'email']
notifier.all_channels() # 所有支持的渠道名notifier.send_text(title="报告", content="...")
│
▼
ChannelRegistry:从 NotificationConfig 实例化所有已配置的 senders
│
▼
ParallelDispatcher(ThreadPoolExecutor):
│
├─ feishu.send_text(...) ──→ ✓ SendResult(channel="feishu", ok=True)
├─ email.send_text(...) ──→ ✓ SendResult(channel="email", ok=True)
├─ wechat.send_text(...) ──→ ✗ SendResult(channel="wechat", ok=False, error="...")
└─ telegram.send_text(...) ──→ ✓ SendResult(channel="telegram", ok=True)
(任一渠道失败,不影响其他)
│
▼
返回 List[SendResult]
pulsefan 本身只懂"发消息"。"生成股票分析报告 + 决定哪些渠道发完整版/摘要版"这个业务语义放在 stocklens 侧:
# stocklens/notification/orchestrator.py
from pulsefan import Notifier
from stocklens.notification.config_adapter import config_to_notification_config
class ReportOrchestrator:
def __init__(self, config):
self.notifier = Notifier(config_to_notification_config(config))
def send_daily_report(self, results):
# 业务编排:生成 markdown / 摘要 / 飞书卡片
full_md = self._render_full_markdown(results)
summary_md = self._render_summary(results)
# 不同渠道发不同版本
self.notifier.sender("feishu").send_card(self._to_feishu_card(results))
self.notifier.sender("email").send_text(title="...", content=full_md)
self.notifier.send_text(title="盘后摘要", content=summary_md)Config 单一映射点:stocklens/notification/config_adapter.py 是唯一把 stocklens.Config(130+ 环境变量)翻译成 pulsefan.NotificationConfig 的地方。新增渠道字段必须改这个文件 + pulsefan/config/schema.py。
把 quantcore + tickbridge + querybus + pulsefan 4 个独立包"组装"成"股票分析"业务的胶水层。
- ❌ 不发明算法——评分/指标全部委托 quantcore
- ❌ 不直接抓数据——全部走 tickbridge / querybus
- ❌ 不直接发推送——全部走 pulsefan
- ✅ 业务编排:流水线、Agent、报告生成、配置管理
- ✅ 持久化:SQLAlchemy ORM、9 个 Repository
- ✅ 领域薄壳:把通用包包装成股票领域 API
stocklens/
│
├── pipeline/ ★ 核心:分析流水线编排
│ ├── orchestrator.py StockAnalysisPipeline(主入口)
│ ├── stages/ 9 个原子 stage(realtime / chip / fundamental /
│ │ intel / enrichment / scoring / metadata / llm 等)
│ ├── data_collector.py DataCollector — 数据采集 + 增量更新 + 多周期 K
│ ├── stock_indicators.py trend + tech_indicators 共用 helper
│ ├── us_data_enricher.py US/HK 并行增强
│ ├── analysis_flow.py 上下文增强 + Agent 路径
│ ├── push_handler.py 报告保存 + 推送调度
│ ├── trading_calendar.py 交易日历(exchange-calendars,fail-open)
│ └── metrics.py pipeline 计时与统计
│
├── agent/ ★ 多 Agent 系统(30 文件)
│ ├── orchestrator.py AgentOrchestrator — multi 模式编排 5 Agent
│ ├── runner.py AgentRunner — ReAct 循环
│ ├── factory.py Agent 工厂
│ ├── agents/ 5 个 Agent:
│ │ technical / intel / risk(可一票否决) /
│ │ portfolio / decision
│ ├── tools/ ToolRegistry + market/data/search/analysis/backtest 工具
│ ├── strategies/ StrategyRouter / StrategyAggregator / StrategyAgent
│ └── skills/ YAML 策略加载器
│
├── analyzer/ ★ LLM 分析器
│ ├── llm_analyzer.py GeminiAnalyzer — LiteLLM Router 调 LLM
│ ├── prompt_builder.py PromptBuilder — system+user prompt 组装
│ ├── response_parser.py JSON 提取 + json_repair
│ ├── result_types.py AnalysisResult dataclass + 8 个一致性校准函数
│ └── trend_analyzer.py 薄壳 → quantcore.indicators.trend
│
├── services/ ★ 业务服务(30 文件)
│ ├── scoring/ 薄壳 → quantcore.scoring(保留 CBOE fetcher 注入)
│ ├── bollinger_service.py 薄壳 → quantcore.indicators.bollinger(保留 IO + 缓存)
│ ├── tech_indicators_service.py 薄壳 → quantcore.stats.correlation
│ ├── task_queue.py AnalysisTaskQueue — 异步任务队列
│ ├── system_config_service.py 运行时配置(.env 读写)
│ ├── history_service.py 分析历史
│ ├── finviz_service.py Finviz 数据领域包装
│ └── ...
│
├── notification/ ★ 通知壳(包装 pulsefan.Notifier)
│ ├── orchestrator.py ReportOrchestrator — daily/batch/metadata 报告编排
│ ├── config_adapter.py Config → pulsefan.NotificationConfig 单一映射
│ ├── metadata_report.py 元数据报告入口
│ └── metadata_sections/ 12 个语义段子包
│
├── search/ ★ 搜索壳(包装 querybus.SearchProvider)
│ └── stock_search_service.py StockSearchService + 5 维度并发情报
│
├── market/ ★ 美股市场环境(无 stocklens 内部依赖,零外部业务依赖)
│ ├── us_macro.py USMacroSnapshot — 12 宏观指标
│ ├── us_market_env.py 四大指数 + PE 百分位
│ ├── profile.py 市场画像
│ └── strategy.py 策略匹配
│
├── portfolio/ ★ 组合管理
│ ├── portfolio_service.py FIFO/均价成本 + FX 换算 + 快照回放
│ ├── import_service.py 雪球/富途/通用格式导入
│ └── risk_service.py 组合风险评估
│
├── backtest/ ★ 回测调度(engine 已下沉到 quantcore)
│ ├── engine.py 薄壳 → quantcore.backtest.engine
│ └── backtest_service.py BacktestService — IO + DB + 编排
│
├── storage/ SQLAlchemy ORM
│ ├── engine.py DatabaseManager 单例 + PRAGMA(WAL/mmap)
│ ├── models.py ORM 表(含 stock_bars 多周期 + stock_quote_snapshot)
│ └── helpers.py scoring 快照保存等辅助
│
├── repositories/ 9 个 Repository(CRUD + 批量 upsert)
│ ├── analysis_repo.py
│ ├── stock_repo.py 日线/周线/多周期 K
│ ├── news_repo.py
│ ├── fundamental_repo.py
│ ├── macro_repo.py
│ ├── backtest_repo.py
│ ├── conversation_repo.py
│ ├── llm_usage_repo.py
│ └── portfolio_repo.py
│
├── config/ Config 单例(130+ 环境变量)
│ ├── settings.py Config 类
│ ├── settings_helpers.py LLM/News/Agent 工具函数
│ ├── llm.py LLM_CHANNELS → LiteLLM Router 解析
│ └── fields/ 字段注册表(按 category 拆分)
│
├── contracts/ ★ 跨层共享类型(核心层 ⇄ 适配层共用)
│ └── bot_message.py BotMessage / ChatType
│
├── assets/ AssetManager — 集中管理 catalog/icons/logos
├── utils/ 并发池 / 二级缓存 / 日志 / markdown→图片
├── auth/ 可选认证
├── schemas/ 报告 schema
└── data/ 数据目录
# stocklens/pipeline/orchestrator.py
class StockAnalysisPipeline:
def __init__(self, config, db, fetcher_manager, ...):
# DI 各种依赖
self.data_collector = DataCollector(...)
self.scoring_service = ScoringService() # quantcore 包装
self.notifier = ReportOrchestrator(config) # pulsefan 包装
self.search_service = StockSearchService() # querybus 包装
def run(self, stock_list):
# Phase 1: 共享数据预取(PREFETCH 池)
shared = self._prefetch_shared_context()
# Phase 2: 单股数据采集(DATA_FETCH 池)
for code in stock_list:
self.data_collector.fetch_and_save(code)
# Phase 3: enrich + LLM 双池并行(PIPELINE 池)
with ThreadPoolExecutor() as enrich_pool, ThreadPoolExecutor() as llm_pool:
for code in stock_list:
future = enrich_pool.submit(self._enrich_one, code, shared)
# 关键加速:enrich 完成立即提交 LLM
future.add_done_callback(
lambda f, c=code: llm_pool.submit(self._llm_analyze, c, f.result())
)
# 报告 + 推送
self.push_handler.dispatch(results)| stocklens 子目录 | 主要 import 的子包 | 起的作用 |
|---|---|---|
pipeline/ |
tickbridge, quantcore, querybus |
编排数据采集 → 算分 → 调情报 |
agent/ |
quantcore, querybus |
Agent 工具调用 quantcore 算分;调 querybus 搜情报 |
analyzer/ |
quantcore.indicators.trend(仅趋势 helper) |
LLM 分析时附带趋势数据 |
services/scoring/ |
quantcore.scoring |
薄壳 + CBOE 期权抓取注入 |
services/bollinger_service.py |
quantcore.indicators.bollinger |
薄壳 + akshare/雪球抓取 + 二级缓存 |
services/tech_indicators_service.py |
quantcore.stats.correlation |
薄壳 + SPY/QQQ 抓取 + 4h 缓存 |
notification/ |
pulsefan |
业务报告编排 + Config 映射 |
search/ |
querybus |
5 维度并发情报 + 股票领域 query 拼接 |
backtest/ |
quantcore.backtest |
薄壳 + IO/DB 调度 |
market/ |
tickbridge(间接) |
12 宏观指标 + 大盘环境 |
portfolio/ |
tickbridge.codes |
FIFO/均价成本 + 行情联动 |
api/
├── app.py create_app() 工厂
│ CORS / 中间件 / lifespan(long_pool/io_pool)
├── deps.py 依赖注入
│
├── middlewares/ 认证 / 限流 / 日志
│
└── v1/endpoints/ ★ 10 个 endpoint 模块
├── analysis.py POST /api/v1/analysis 触发分析
├── agent.py Agent 对话
├── history.py 历史查询
├── portfolio.py 组合管理
├── backtest.py 回测调度
├── stocks.py 股票元数据
├── sectors.py 板块数据
├── system_config.py 运行时配置(.env 读写)
├── auth.py 登录
└── usage.py LLM 用量统计
双线程池架构:
# api/app.py lifespan
@asynccontextmanager
async def lifespan(app):
app.state.long_pool = ThreadPoolExecutor(max_workers=N) # LLM/Agent(重)
app.state.io_pool = ThreadPoolExecutor(max_workers=M) # 轻 IO
yield
app.state.long_pool.shutdown()
app.state.io_pool.shutdown()LLM 调用走 long_pool(避免占满轻 IO 通道),普通查询走 io_pool。
bot/
├── platforms/ 平台 adapter
│ ├── dingtalk_stream.py 钉钉 Stream
│ └── feishu_stream.py 飞书 Stream
├── commands/ ★ 6 个命令处理器
│ ├── help.py /help
│ ├── analyze.py /analyze AAPL
│ ├── batch.py /batch
│ ├── ask.py /ask
│ ├── chat.py /chat
│ └── status.py /status
├── dispatcher.py CommandDispatcher — 统一命令分发
└── models.py re-export stocklens.contracts.BotMessage
关键设计:bot/ 永不直接 import stocklens 内部模块。所有跨层共享类型(BotMessage, ChatType)都通过 stocklens.contracts/ 暴露,bot/ 从那里读。
用户在 Web UI 点击"分析 AAPL"
│
▼
api/v1/endpoints/analysis.py ← 适配层
│ POST /api/v1/analysis {"codes": ["AAPL"]}
│ 提交 Future 到 app.state.long_pool
▼
stocklens/pipeline/orchestrator.py ← 业务编排
│ StockAnalysisPipeline.run(["AAPL"])
│
├──▶ Phase 1: 共享数据预取
│ │
│ ├─▶ stocklens/market/us_macro.py
│ │ └─▶ tickbridge.DataProvider.get_realtime_quote(".VIX") ← 数据层
│ │ └─▶ tickbridge/sources/yfinance.py ← 实际抓取
│ │
│ └─▶ stocklens/market/us_market_env.py
│ └─▶ tickbridge.DataProvider.get_main_indices("us")
│
├──▶ Phase 2: 单股数据采集
│ │
│ └─▶ stocklens/pipeline/data_collector.py
│ ├─▶ tickbridge.DataProvider.get_realtime_quote("AAPL")
│ ├─▶ tickbridge.DataProvider.get_daily_data("AAPL", days=200)
│ ├─▶ tickbridge.DataProvider.get_fundamental("AAPL")
│ ├─▶ tickbridge.DataProvider.finviz_forward("AAPL") ← 域工具
│ └─▶ stocklens/repositories/stock_repo.py 保存到 SQLite
│
└──▶ Phase 3: enrich + LLM 双池并行
│
├──▶ enrich_pool: stocklens/pipeline/stages/enrichment_stage.py
│ │
│ ├─▶ quantcore.indicators.StockTrendAnalyzer.analyze(df) ← 算法层
│ │ 返回 TrendSnapshot
│ │
│ ├─▶ stocklens/services/bollinger_service.py
│ │ ├─ akshare 抓日线 + 雪球抓 60m(IO)
│ │ └─▶ quantcore.indicators.compute_bollinger(df) ← 算法层
│ │
│ ├─▶ stocklens/services/tech_indicators_service.py
│ │ ├─ akshare 抓 SPY/QQQ(IO)
│ │ └─▶ quantcore.indicators.compute_tech_indicators ← 算法层
│ │ quantcore.stats.compute_correlation_pair
│ │
│ ├─▶ stocklens/search/stock_search_service.py
│ │ └─▶ querybus.SearchProvider.search(...) ← 搜索层
│ │ └─▶ querybus/providers/tavily.py ← 实际搜索
│ │
│ ├─▶ stocklens/services/scoring/__init__.py
│ │ └─▶ quantcore.scoring.ScoringEngine.evaluate(ctx) ← 算法层
│ │ 跑 18 个模型,返回 ScoringOverview
│ │
│ └─▶ enrichment 完成 → callback 提交 LLM
│
└──▶ llm_pool: stocklens/pipeline/stages/llm_stage.py
│
├─▶ stocklens/analyzer/prompt_builder.py 拼提示词
│
├─▶ stocklens/analyzer/llm_analyzer.py
│ └─ LiteLLM Router 调 Claude/GPT/Gemini
│
├─▶ stocklens/analyzer/response_parser.py 解析 JSON
│ 返回 AnalysisResult
│
└─▶ stocklens/analyzer/result_types.py
├─ enforce_stop_loss_consistency
├─ enforce_trend_status_consistency ← 用 quantcore 的 golden_cross 校准
├─ enforce_score_band_consistency
└─ fill_chip_structure_if_needed
│
▼
stocklens/pipeline/push_handler.py ← 推送编排
│
└─▶ stocklens/notification/orchestrator.py
│
├─ stocklens/notification/config_adapter.py 翻译 Config → NotificationConfig
│
└─▶ pulsefan.Notifier.send_text(...) ← 通知层
│
└─▶ ParallelDispatcher fan-out:
├─▶ pulsefan/senders/feishu.py ✓
├─▶ pulsefan/senders/wechat.py ✓
├─▶ pulsefan/senders/email.py ✓
└─▶ pulsefan/senders/telegram.py ✓
(所有渠道并行,互不阻塞)
│
▼
返回 results 给前端
│
▼
[T+10 自动回测]
stocklens/backtest/backtest_service.py
└─▶ quantcore.backtest.BacktestEngine.evaluate_single(...) ← 算法层
评估方向准确率/胜率/止损止盈
| 规则 | 例子 |
|---|---|
| 业务层不发明算法 | 写新评分模型 → 必须放 quantcore/scoring/ |
| 业务层不直接抓数据 | 调雪球 → 必须经过 tickbridge.DataProvider |
| 业务层不直接发推送 | 发飞书消息 → 必须经过 pulsefan.Notifier |
| 业务层包装通用包 | "搜某股新闻" → 在 stocklens.search/ 包装 querybus |
| 算法层不接收外部数据 | quantcore 函数都是 pandas DataFrame in / dataclass out |
| 跨层共享类型 | 放在 stocklens.contracts/(如 BotMessage) |
┌──────────────────────────────────────────────────────────────────────┐
│ │
│ 触发:Web UI / REST / Bot │
│ │ │
│ ▼ │
│ api / bot │
│ │ 转换成 BotMessage(stocklens.contracts) │
│ ▼ │
│ stocklens.pipeline.StockAnalysisPipeline │
│ │ │
│ ├─读数据────────▶ tickbridge.DataProvider │
│ │ │ ├─缓存 / SQLite / 10 sources failover │
│ │ │ ▼ │
│ │ │ 上游 API(雪球/Yahoo/Finviz/...) │
│ │ ▼ │
│ │ pandas DataFrame │
│ │ │
│ ├─搜情报────────▶ querybus.SearchProvider │
│ │ │ ├─缓存 / 7 providers fall-over │
│ │ │ ▼ │
│ │ │ Tavily/Bocha/Brave/... │
│ │ ▼ │
│ │ SearchResponse │
│ │ │
│ ├─跑算法────────▶ quantcore.ScoringEngine.evaluate(ctx) │
│ │ │ ├─18 模型纯计算 │
│ │ │ ▼ │
│ │ │ ScoringOverview(dataclass) │
│ │ ▼ │
│ │ composite_score(0-100) │
│ │ │
│ ├─调 LLM────────▶ analyzer.LiteLLM Router │
│ │ │ ▼ │
│ │ │ Claude/GPT/Gemini │
│ │ ▼ │
│ │ AnalysisResult(被 quantcore 数据校准) │
│ │ │
│ ├─入库────────▶ storage/repositories │
│ │ │ ▼ │
│ │ │ SQLite (WAL mode) │
│ │ ▼ │
│ │ AnalysisRecord │
│ │ │
│ ├─推送────────▶ pulsefan.Notifier.send_text │
│ │ │ ├─10 senders 并行扇出 │
│ │ │ ▼ │
│ │ │ 飞书/微信/邮件/Telegram/... │
│ │ ▼ │
│ │ List[SendResult] │
│ │ │
│ └─T+N 回测─▶ quantcore.BacktestEngine.evaluate_single │
│ │ │
│ ▼ │
│ 方向准确率 / 胜率 / 收益率 │
│ │
└──────────────────────────────────────────────────────────────────────┘
| 包 | 入口文件 | 唯一对外类 |
|---|---|---|
| quantcore | quantcore/__init__.py |
ScoringEngine, BacktestEngine |
| tickbridge | tickbridge/core/provider.py |
DataProvider |
| querybus | querybus/core/provider.py |
SearchProvider |
| pulsefan | pulsefan/core/notifier.py |
Notifier |
| stocklens | stocklens/pipeline/orchestrator.py |
StockAnalysisPipeline |
| 文件 | 职责 |
|---|---|
main.py |
启动 FastAPI + 可选 Bot Stream |
api/app.py |
create_app() 工厂;lifespan 创建 long_pool/io_pool |
stocklens/config/settings.py |
Config 单例,130+ 环境变量 |
stocklens/pipeline/orchestrator.py |
流水线主入口 |
stocklens/notification/config_adapter.py |
Config → NotificationConfig 单一映射 |
stocklens/services/scoring/__init__.py |
ScoringService = ScoringEngine + CBOE fetcher |
| 套件 | 路径 | 用例数 |
|---|---|---|
| 主套件 | tests/unit/ |
341 |
| querybus | querybus/tests/ |
31 |
| pulsefan | pulsefan/tests/ |
19 |
| tickbridge | tickbridge/tests/ |
多个文件 |
运行命令:
pytest tests/ -q # 主项目
pytest querybus/tests/ -q # querybus 独立
pytest pulsefan/tests/ -q # pulsefan 独立
pytest tickbridge/tests/ -q # tickbridge 独立文档维护:本文档随
docs/CONTEXT.md与CLAUDE.md同步更新。 模块层面更细的类/函数清单见docs/modules/<模块名>.md。