
数据流图源码
flowchart TB
%% ═══════════════════════════════════════════════════
%% 外部数据源
%% ═══════════════════════════════════════════════════
subgraph DataSources["📡 外部数据源"]
TDX_STD["easy_tdx TdxClient<br/>(通达信标准协议)"]
TDX_MAC["easy_tdx MacClient<br/>(通达信Mac协议)"]
TDX_MACEX["easy_tdx MacExClient<br/>(通达信Mac扩展协议)"]
EMoney["东方财富<br/>(Playwright)"]
MySQL_DB[("MySQL<br/>stock_pool 数据库")]
end
%% ═══════════════════════════════════════════════════
%% 采集层
%% ═══════════════════════════════════════════════════
subgraph Collectors["⚙️ 数据采集层"]
direction TB
TC["TickCollector<br/>逐笔成交 @15s"]
OBC["OrderBookCollector<br/>五档盘口 @10s"]
LUC["LimitUpCollector<br/>涨停检测 @60s"]
HKC["HKCollector<br/>港股行情+Tick @30s"]
SDF["SectorDataFetcher<br/>30个一级行业排行 @30s"]
MDP["MarketDataProvider<br/>流通市值/换手率 @60s"]
LHC["LonghuCollector<br/>龙虎榜 @17:00"]
SPL["StockPoolLoader<br/>股票池热更新 @5min"]
end
%% ═══════════════════════════════════════════════════
%% 内存消息总线
%% ═══════════════════════════════════════════════════
subgraph Bus["🔀 DataBus (asyncio.Queue)"]
direction LR
CH_TICK[/"tick 通道"/]
CH_OB[/"orderbook 通道"/]
CH_LU[/"limitup 通道"/]
end
%% ═══════════════════════════════════════════════════
%% 计算引擎层
%% ═══════════════════════════════════════════════════
subgraph Engines["🧠 计算引擎层"]
direction TB
FE["FeatureEngine<br/>25维特征计算 @60s"]
HSE["HotScoreEngine<br/>热点评分+排行"]
SE["SectorEngine<br/>板块联动(混合模式)"]
AE["AlertEngine<br/>异动告警"]
end
%% ═══════════════════════════════════════════════════
%% 评估层
%% ═══════════════════════════════════════════════════
subgraph Eval["📊 评估层"]
direction TB
EE["EvalEngine<br/>快照→回填→Precision"]
DR["DailyReporter<br/>日报+权重调优"]
end
%% ═══════════════════════════════════════════════════
%% 存储层
%% ═══════════════════════════════════════════════════
subgraph Storage["💾 存储层"]
direction LR
Redis[("Redis<br/>排行/特征/告警<br/>状态持久化")]
DuckDB[("DuckDB<br/>特征历史/评估")]
end
%% ═══════════════════════════════════════════════════
%% API层
%% ═══════════════════════════════════════════════════
subgraph API["🌐 API 服务层 (FastAPI :8001)"]
direction LR
REST["REST API<br/>/hotspot /sector<br/>/alert /eval /longhu"]
WS["WebSocket /ws<br/>实时推送"]
end
%% ═══════════════════════════════════════════════════
%% 数据源 → 采集器 连接
%% ═══════════════════════════════════════════════════
TDX_STD -->|"逐笔/盘口/涨停<br/>TCP连接池x8"| TC
TDX_STD --> OBC
TDX_STD --> LUC
TDX_MAC -->|"get_board_ranking<br/>YJ_LEVEL1"| SDF
TDX_MAC -->|"get_stock_quotes<br/>float_shares"| MDP
TDX_MACEX -->|"goods_quotes<br/>goods_transaction"| HKC
TDX_MACEX -->|"goods_quotes<br/>港股市值/换手"| MDP
EMoney -->|"龙虎榜JSON"| LHC
MySQL_DB -->|"stock_pool<br/>symbol/sector"| SPL
%% ═══════════════════════════════════════════════════
%% 采集器 → DataBus
%% ═══════════════════════════════════════════════════
TC -->|"TickRecord[]"| CH_TICK
HKC -->|"TickRecord[]"| CH_TICK
OBC -->|"OrderBookSnapshot[]"| CH_OB
HKC -->|"OrderBookSnapshot[]"| CH_OB
LUC -->|"LimitUpRecord[]"| CH_LU
%% ═══════════════════════════════════════════════════
%% 采集器 → 引擎(直接注入,不经过DataBus)
%% ═══════════════════════════════════════════════════
SDF -.->|"BoardRankingItem[]<br/>change_pct/up_ratio/net_buy"| SE
MDP -.->|"MarketDataItem<br/>float_cap/turnover"| FE
SPL -.->|"sector_mapping<br/>watch_list更新"| SE
%% ═══════════════════════════════════════════════════
%% DataBus → 引擎
%% ═══════════════════════════════════════════════════
CH_TICK --> FE
CH_OB --> FE
CH_LU --> SE
%% ═══════════════════════════════════════════════════
%% 引擎间交互
%% ═══════════════════════════════════════════════════
FE -->|"FeatureSnapshot{25维}"| HSE
SE -->|"板块联动特征注入<br/>sector_rank/up_ratio"| FE
HSE -->|"HotScore结果"| AE
SE -->|"板块热度"| AE
HSE -->|"Top N 快照"| EE
%% ═══════════════════════════════════════════════════
%% 引擎 → 存储
%% ═══════════════════════════════════════════════════
HSE -->|"排行/分数历史"| Redis
FE -->|"特征快照"| Redis
FE -->|"特征历史"| DuckDB
SE -->|"板块排行"| Redis
AE -->|"告警列表"| Redis
EE -->|"评估记录"| DuckDB
LHC -->|"龙虎榜记录"| Redis
DR -->|"日报JSON"| DuckDB
%% ═══════════════════════════════════════════════════
%% 存储 → API
%% ═══════════════════════════════════════════════════
Redis --> REST
DuckDB --> REST
HSE -->|"ranking_update"| WS
AE -->|"alert事件"| WS
EE -->|"eval_update"| WS
%% ═══════════════════════════════════════════════════
%% 优雅停机/暖启动
%% ═══════════════════════════════════════════════════
FE <-->|"minute_agg<br/>持久化/恢复"| Redis
TC <-->|"last_index<br/>持久化/恢复"| Redis
HSE <-->|"score_history<br/>持久化/恢复"| Redis
%% ═══════════════════════════════════════════════════
%% 样式
%% ═══════════════════════════════════════════════════
classDef source fill:#e1f5fe,stroke:#0288d1
classDef collector fill:#fff3e0,stroke:#f57c00
classDef bus fill:#f3e5f5,stroke:#7b1fa2
classDef engine fill:#e8f5e9,stroke:#388e3c
classDef eval fill:#fce4ec,stroke:#c62828
classDef storage fill:#fff9c4,stroke:#f9a825
classDef api fill:#e8eaf6,stroke:#3f51b5
class TDX_STD,TDX_MAC,TDX_MACEX,EMoney,MySQL_DB source
class TC,OBC,LUC,HKC,SDF,MDP,LHC,SPL collector
class CH_TICK,CH_OB,CH_LU bus
class FE,HSE,SE,AE engine
class EE,DR eval
class Redis,DuckDB storage
class REST,WS api

采集器数据流图源码
sequenceDiagram
participant DS as easy_tdx 数据源
participant TC as TickCollector
participant OBC as OrderBookCollector
participant SDF as SectorDataFetcher
participant MDP as MarketDataProvider
participant DB as DataBus
participant FE as FeatureEngine
participant SE as SectorEngine
participant HSE as HotScoreEngine
participant AE as AlertEngine
participant R as Redis
participant WS as WebSocket
Note over DS,WS: ═══ 每分钟计算周期 ═══
rect rgb(255, 243, 224)
Note right of DS: 采集阶段 (每15-30s)
DS->>TC: get_transaction_data (A股逐笔)
DS->>OBC: get_security_quotes (五档盘口)
DS->>SDF: get_board_ranking (30个板块)
DS->>MDP: get_stock_quotes (市值/换手)
TC->>DB: publish(tick, TickRecord[])
OBC->>DB: publish(orderbook, Snapshot[])
SDF-->>SE: 直接注入 BoardRankingItem[]
MDP-->>FE: 直接注入 MarketDataItem[]
end
rect rgb(232, 245, 233)
Note right of FE: 特征计算阶段 (@60s)
DB->>FE: consume(tick) + consume(orderbook)
FE->>FE: 计算25维特征(量比/主买/盘口...)
SE->>FE: 注入板块联动特征(rank/up_ratio)
FE->>FE: 合并 float_cap + turnover_pct
end
rect rgb(200, 230, 255)
Note right of HSE: 评分阶段
FE->>HSE: FeatureSnapshot[] (全股票池)
SE->>SE: _compute_hybrid_metrics<br/>(板块级change_pct + 逐笔surge/limitup)
HSE->>HSE: 加权评分 → 排行 → 等级划分
HSE->>AE: 检测异动(score急升/量爆发/板块共振)
end
rect rgb(252, 228, 236)
Note right of R: 输出阶段
HSE->>R: ZADD hotspot:ranking (Top N)
SE->>R: ZADD sector:ranking (板块排行)
AE->>R: LPUSH alert:recent
HSE->>WS: push ranking_update
AE->>WS: push alert
end
暂无评论,快来抢沙发吧~