Risk-MAS 框架图说明
整体流程
输入 → 1.状态模块 → 2.输入规范化 → 3.数据与指标模块 → 4.编排模块 → 5/6.并行分析 → 7.汇总与决策 → 输出
架构设计原则
1. 职责分离(Separation of Concerns)
- Gatekeeper:数据可用性检查 → 回答"哪些节点可以运行"
- Supervisor:业务逻辑选择 → 回答"哪些节点需要运行"
- 避免重复判断,单一数据源
2. 依赖注入(Dependency Injection)
RuntimeConfig对象注入到关键模块(supervisor, agents)- 显式依赖,易于测试和维护
- 避免隐式的
os.getenv()调用
3. 显式状态管理(Explicit State)
- 使用
TypedDict定义RiskState - 类型安全,可追溯
- 所有节点通过 state 传递数据
4. 并行执行(Parallel Execution)
- 使用 LangGraph
Send API实现真正的并行 - 6 个分析节点并发执行
- 提升处理效率
模块详细说明
1. 状态模块
┌─────────────────┐
│ state.py │
│ (显式State定义) │
└─────────────────┘
2. 输入规范化模块
┌─────────────────┐
│ validate.py │
│ (校验与归一化) │
└─────────────────┘
3. 数据与指标模块
┌─────────────────┬─────────────────┬─────────────────┐
│ csv_data.py │ data_quality.py │ snapshot.py │
│ (数据读取) │ (数据质量检查) │ (指标快照) │
└─────────────────┴─────────────────┴─────────────────┘
│
▼
┌─────────────────────┐
│ utils.py │
│ (共享工具函数) │
│ normalize_weights │
│ compute_hhi │
│ compute_effective_n │
└─────────────────────┘
4. 编排模块
┌─────────────────────────────────────────────────────┐
│ graph.py │
│ (流程编排与并行调度) │
└─────────────────────────────────────────────────────┘
│
▼
┌─────────────────┐ │
│ gatekeeper.py │───────┤
│ (数据可用性检查) │ │
│ (候选节点筛选) │ │
└─────────────────┘ │
│
▼
┌─────────────────┐
│ supervisor.py │
│ (业务逻辑选择) │
└─────────────────┘
│
│ Send API (并行分发)
▼
5. 分析链路模块 + 6. Agent模块(并行执行)
┌─────────────────────────────────────────┐
│ 并行执行区域 │
│ ┌─────────────┐ ┌─────────────┐ │
│ │ market.py │ │concentration│ │
┌─────────┼─→│ (市场风险) │ │ .py │←──────┼─┐
│ │ └─────────────┘ │ (集中度) │ │ │
│ │ └─────────────┘ │ │
supervisor├─────────┼─→┌─────────────┐ ┌─────────────┐←──────┼─┤
│ │ │diversific- │ │liquidity.py │ │ │
│ │ │ ation.py │ │ (流动性) │ │ │
│ │ │ (分散度) │ └─────────────┘ │ │
│ │ └─────────────┘ │ │
│ │ │ │
│ │ ┌─────────────┐ ┌─────────────┐ │ │
└─────────┼─→│macro_agent │ │compliance_ │←──────┼─┘
│ │ .py │ │ agent.py │ │
│ │ 宏观Agent │ │ 合规Agent │ │
│ └─────────────┘ └─────────────┘ │
│ │
└─────────────────────────────────────────┘
│
│ 全部汇聚
▼
7. 汇总与决策模块
┌─────────────────┐
│ reducer.py │
│ (风险汇总输出) │
└────────┬────────┘
│
▼
┌─────────────────┐ ┌─────────────────┐
│ decision.py │────→│ solver.py │ (restrict时触发)
│ (决策引擎) │ │ (约束求解) │
└────────┬────────┘ └─────────────────┘
│
├── pass ──→ 放行
├── warn ──→ 预警
├── restrict → 限制 + 调仓建议
└── block ──→ 阻断
│
▼
┌─────────────────┐
│ audit.py │
│ (审计与追溯) │
└─────────────────┘
│
▼
完成
Skills 模块
┌─────────────────────────────────────────────────────┐
│ skills_runtime.py │
│ (读取SKILL/Schema/工具白名单) │
└─────────────────────────────────────────────────────┘
│
▼
┌────────────────────────────────────────────────────────────────┐
│ skills/ │
│ ├── macro-tool-calling/ │
│ │ ├── SKILL.md (技能定义) │
│ │ └── output.schema.json (输出结构) │
│ ├── compliance-tool-calling/ │
│ │ ├── SKILL.md │
│ │ └── output.schema.json │
│ ├── snippets/ │
│ │ ├── evidence_rules.md (证据规范) │
│ │ └── decision_rubric.md (决策标准) │
│ └── tools/ │
│ └── tool_interfaces.yaml (工具注册表) │
└────────────────────────────────────────────────────────────────┘
配置模块
┌─────────────────────────────────────────────────────────────────┐
│ config.py │
│ RuntimeConfig (dataclass) │
├─────────────────────────────────────────────────────────────────┤
│ 特性: │
│ • @dataclass(frozen=True) - 不可变配置对象 │
│ • from_env() - 统一读取环境变量 │
│ • 类型安全辅助函数(_env_int / _env_float / _env_bool) │
│ • 向后兼容别名:Config = RuntimeConfig │
├─────────────────────────────────────────────────────────────────┤
│ 配置分类: │
│ 📊 数据与采样: │
│ • csv_data_dir / sample_universe_size / random_seed │
│ • market_lookback_days / asof_date │
│ 🌍 宏观数据: │
│ • tushare_token / macro_series_config │
│ • macro_stale_days / macro_severity_weight │
│ 📋 合规与 RAG: │
│ • compliance_rag_source / rag_engine │
│ 🤖 LLM 配置: │
│ • llm_model / openai_api_key / openai_base_url │
│ • enable_supervisor │
│ 💼 组合与求解器: │
│ • default_aum / target_holdings / cash_symbol │
│ • lp_turnover_weight / lp_solver │
├─────────────────────────────────────────────────────────────────┤
│ 使用方式: │
│ # 创建配置实例 │
│ config = RuntimeConfig.from_env() │
│ # 依赖注入到关键模块 │
│ supervisor_chain(state, llm, candidates, config) │
│ run_macro_agent(state, llm, config) │
│ constraint_solver(state, config) │
└─────────────────────────────────────────────────────────────────┘
Tools 模块汇总
┌──────────────────────────────────────────────────────────────┐
│ tools/ │
├──────────────────────────────────────────────────────────────┤
│ utils.py 共享工具函数 │
│ validate.py 输入校验与归一化 │
│ csv_data.py CSV数据读取 │
│ data_quality.py 数据质量检查 │
│ snapshot.py 指标快照计算 │
│ constraints.py 规则评估与硬约束 │
│ decision.py 决策引擎 │
│ solver.py 约束求解(cvxpy) │
│ rules.py 规则加载 │
│ audit.py 审计输出 │
│ calibrate_rules.py 规则阈值校准 │
│ calibrate_macro_series.py 宏观序列阈值校准 │
└──────────────────────────────────────────────────────────────┘
数据流图
State 流转概览
┌─────────────────────────────────────────────────────────────────┐
│ RiskState (TypedDict) │
│ 贯穿整个工作流的显式状态 │
└─────────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ INPUT: intent + context │
│ ├─ intent: {date, mode, targets} │
│ └─ context: {current_positions, policy_profile, aum, ...} │
└────────────────────────┬────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ VALIDATE: 校验与归一化 │
│ OUTPUT: normalized (归一化后的输入) │
└────────────────────────┬────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ DATA_QUALITY: 数据质量检查 │
│ OUTPUT: data_quality {market, macro, compliance, positions} │
└────────────────────────┬────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ SNAPSHOT: 指标快照计算 │
│ OUTPUT: snapshot_metrics {portfolio_volatility, hhi, ...} │
└────────────────────────┬────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ GATEKEEPER: 数据可用性检查 │
│ OUTPUT: candidate_nodes (可运行的节点列表) │
└────────────────────────┬────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ SUPERVISOR: 业务逻辑选择 │
│ INPUT: candidate_nodes, snapshot_metrics, rule_findings │
│ OUTPUT: nodes_to_run, pending_agents (需要运行的节点) │
└────────────────────────┬────────────────────────────────────────┘
│
▼ (Send API 并行分发)
┌─────────────────────────────────────────────────────────────────┐
│ PARALLEL ANALYSIS: 6个节点并发执行 │
│ ├─ market → finding_market │
│ ├─ concentration → finding_concentration │
│ ├─ diversification → finding_diversification │
│ ├─ liquidity → finding_liquidity │
│ ├─ macro → finding_macro + tool_calls_macro │
│ └─ compliance → finding_compliance + tool_calls_compliance │
└────────────────────────┬────────────────────────────────────────┘
│ (全部汇聚)
▼
┌─────────────────────────────────────────────────────────────────┐
│ REDUCER: 汇总所有 findings │
│ OUTPUT: risk_summary (汇总的风险报告) │
└────────────────────────┬────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ CONSTRAINTS: 约束评估 │
│ OUTPUT: rule_findings, binding_constraints │
└────────────────────────┬────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ DECISION: 决策引擎 │
│ OUTPUT: decision {decision, recommendation, ...} │
└────────────────────────┬────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ SOLVER: 约束求解 (仅在 restrict 时触发) │
│ OUTPUT: recommended_actions (调仓建议) │
└────────────────────────┬────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ AUDIT: 审计日志 │
│ OUTPUT: audit (完整的审计追踪) │
└────────────────────────┬────────────────────────────────────────┘
│
▼
最终输出
模块间交互说明
1. 输入层 → 预处理管道
Validate (输入验证)
- 输入:
intent,context - 输出:
normalized,validation - 作用: 校验输入格式,归一化权重,计算 asof_date
- 依赖: 无
Data Quality (数据质量检查)
- 输入:
normalized - 输出:
data_quality - 作用: 检查市场数据、宏观数据、合规数据的可用性和新鲜度
- 依赖: CSV 数据文件、Tushare API(可选)
Snapshot (指标快照)
- 输入:
normalized,data_quality - 输出:
snapshot_metrics - 作用: 计算组合指标(波动率、HHI、流动性等)
- 依赖: 市场数据
2. 编排层
Gatekeeper (门控)
- 输入:
data_quality - 输出:
candidate_nodes - 作用: 基于数据可用性筛选候选节点
- 规则:
- 如果
timeseries_available=False,则 macro 不进入候选 - 如果
text_available=False,则 compliance 不进入候选 - 其他节点默认进入候选
- 如果
Supervisor (调度器)
- 输入:
candidate_nodes,snapshot_metrics,rule_findings,data_quality - 输出:
nodes_to_run,pending_agents,supervisor_rationale - 作用: 基于业务逻辑从候选节点中选择需要运行的节点
- 依赖: LLM(可选,不可用时回退到运行所有候选节点)
- 注入:
RuntimeConfig对象
3. 并行分析层
确定性链路节点 (market, concentration, diversification, liquidity)
- 输入:
snapshot_metrics,normalized - 输出:
finding_<name>(包含 severity, summary, metrics) - 执行条件: 在
pending_agents列表中 - 特点: 无 LLM 调用,纯计算逻辑
Agent 节点 (macro, compliance)
- 输入:
snapshot_metrics,normalized,data_quality - 输出:
finding_<name>,tool_calls_<name> - 执行条件: 在
pending_agents列表中 - 特点: 使用 LLM + 工具调用
- 依赖: LLM, Tushare API (macro), RAG (compliance)
- 注入:
RuntimeConfig对象
并行机制:
- 使用 LangGraph
Send API - Supervisor 返回
List[Send] - 所有节点真正并发执行
- 全部完成后汇聚到 Reducer
4. 后处理管道
Reducer (汇总)
- 输入: 所有
finding_*字段 - 输出:
risk_summary - 作用: 汇总所有风险发现,计算最高 severity
Constraints (约束评估)
- 输入:
snapshot_metrics,normalized - 输出:
rule_findings,binding_constraints - 作用: 评估硬规则约束(HHI、波动率、流动性等)
- 依赖:
rules.yaml配置文件
Decision (决策引擎)
- 输入:
risk_summary,rule_findings,binding_constraints - 输出:
decision(pass/warn/restrict/block) - 作用: 综合风险报告和规则约束,做出最终决策
- 规则:
- 硬规则优先(rule_level)
- 风险报告兜底(report_level)
Solver (约束求解)
- 输入:
normalized,binding_constraints,snapshot_metrics - 输出:
recommended_actions - 作用: 使用 CVXPY 生成满足约束的调仓建议
- 触发条件: 仅在
decision=restrict时运行 - 注入:
RuntimeConfig对象
Audit (审计)
- 输入: 整个
RiskState - 输出:
audit - 作用: 生成完整的审计追踪,包含所有决策依据
- 注入:
RuntimeConfig对象
5. 数据传递方式
显式 State 管理:
- 所有数据通过
RiskState(TypedDict) 传递 - 每个节点返回
Dict[str, Any],自动合并到 State - 类型安全,可追溯
配置注入:
RuntimeConfig对象通过函数参数注入- 避免隐式的
os.getenv()调用 - 易于测试和 mock
并行安全:
- 每个并行节点接收独立的 State 副本
- 节点间无共享状态
- 通过 Reducer 汇聚结果