跳到主要内容

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 汇聚结果