简介:本资源是一份面向Python开发者与金融科技从业者的AI智能体实战指南,聚焦基于Dify与MCP协议构建可执行金融操作的微信端理财助手,解决行情分析→策略生成→自动化推送的全链路落地难题。资源为单文件PDF文档(248KB),完整覆盖金融助手能力架构、MCP协议配置、Dify工作流编排、Flask微信路由服务、RAG研报知识库集成等核心模块,并附带关键代码片段——如yfinance实时股价调用、蒙特卡洛组合风险评估、马科维茨优化引擎及微信模板消息格式化逻辑。内容预览显示其深度结合量化投资实践,包含意图识别规则、工具函数声明、工作流YAML定义及生产级部署要点,兼顾技术严谨性与金融合规性。目前已有262人学习下载,适合具备Python基础、希望掌握AI+金融垂直场景开发的工程师快速复现高并发下的自动化投顾系统。
1. 为什么一个“微信里自动推理财策略”的系统,非得用 Dify + MCP 搭?——不是炫技,是绕不开的工程现实
你有没有遇到过这种场景:客户在微信里问“我手上有 50 万,能买什么?”;理财经理手动查持仓、翻产品库、比对风险等级、套公式算夏普比率,再截图发过去——整个过程 8 分钟,客户已切到下一个群。这不是服务慢,是人脑+人工流程根本扛不住高频、个性化、强合规的金融决策响应需求。而市面上所谓“智能投顾”SaaS,要么嵌在独立 App 里(用户不打开)、要么走 H5(微信里被限流/跳转拦截/无消息触达),更别说策略逻辑硬编码、改个止盈线就得发版。
这个标题里的「金融科技基于 Dify+MCP 的智能金融理财助手开发:微信端自动化组合策略推送系统构建」,说白了就是:用 Dify 做策略逻辑编排与 LLM 编排中枢,用 MCP(Model Control Protocol)做模型调用层的标准化胶水,把策略生成、合规校验、用户画像匹配、微信消息组装、定时/事件触发全链路串起来,最终让策略像天气预报一样,准时、精准、静默地推到用户微信对话框里。它不替代投顾,而是把投顾最耗时的“策略初筛+话术生成+合规复核”环节自动化掉。适合中小券商财富条线、银行私行部、持牌基金销售机构——他们没资源自研大模型中台,但又必须快速上线可审计、可回溯、可灰度的策略触达能力。核心价值不在“AI”,而在“策略可配置、推送可追踪、话术可审计、模型可替换”这十六个字。
2. Dify 是怎么把“理财策略”变成可拖拽、可调试、可上线的?——从知识库到工作流的三层穿透
Dify 在这里不是当个聊天框外壳,而是作为策略即代码(Strategy-as-Code)的低代码执行引擎。它的价值在于把原本散落在 Excel 公式、Python 脚本、风控文档里的策略逻辑,翻译成可版本管理、可 AB 测试、可人工干预的可视化工作流。下面拆解三层落地动作,每一步都对应真实金融场景的刚性需求。
2.1 知识库不是塞 PDF,而是建“策略合规词典”:用向量+关键词双路召回防踩雷
金融领域最怕“幻觉”——模型胡诌一个不存在的基金代码,或把“R3 风险等级”说成“R2”。Dify 知识库必须承担“策略事实核查员”角色。我们不传《公募基金销售管理办法》全文 PDF,而是结构化拆解:
- 产品白名单表:CSV 格式,含
fund_code,fund_name,risk_level,min_holding_days,is_suspended字段,每日从 TA 系统同步; - 合规话术库:JSONL 格式,每行一条:“当用户风险测评结果为 C2,且询问‘保本’时,禁止出现‘本金保障’‘零风险’等表述,应替换为‘历史业绩不预示未来表现,投资需谨慎’”;
- 术语映射表:YAML 格式,定义“年化 4.5%” → “近一年年化收益率 4.5%,过往业绩不构成未来收益承诺”。
提示:Dify 默认的语义检索在金融长尾词上容易失效(比如“T+0 货币基金”和“快赎货币基金”向量距离远)。我们强制开启关键词增强模式(Keyword Enhancement),并在上传时对字段加前缀标签,如
risk_level: R3→#RISK_LEVEL#R3,确保“R3 用户能否买该产品”这类查询必命中。
# 上传产品白名单 CSV 时,用脚本预处理加标签(关键!) awk -F',' 'NR>1 {print $1","$2",#RISK_LEVEL#" $3 ",#HOLDING_DAYS#" $4 ",#SUSPENDED#" $5}' products.csv > products_tagged.csv这段命令给每一列加上#RISK_LEVEL#这类元标签,Dify 在分块时会保留这些符号,后续检索带#RISK_LEVEL#R3的 query 就能 100% 召回,不依赖向量相似度。这是我们在某城商行落地时血泪经验:纯向量召回导致 17% 的高风险产品被错误推荐给保守型用户,加标签后归零。
2.2 工作流不是画流程图,而是定义“策略决策树”:用条件节点卡死合规红线
Dify 工作流(Workflow)在这里是真正的“策略引擎”。我们不用它写“你好,我是理财助手”,而是建一棵动态决策树,每个节点都是硬性业务规则:
- 起点:微信用户发送消息(通过 Webhook 接入);
- 分支 1(用户画像校验):调用内部 API 获取该用户
risk_level,asset_total,last_trade_date,若risk_level == "C1"且消息含“股票”,则直接返回预设话术“根据您的风险测评结果,暂不建议参与股票类投资”; - 分支 2(策略匹配):若用户资产 > 30 万且近 3 月有交易,则触发“股债平衡策略生成”子工作流;
- 分支 3(合规熔断):所有策略输出前,强制过一遍“话术合规检查”节点——输入待发送文本,调用本地部署的规则引擎(Python Flask),返回
{"pass": false, "error": "出现'稳赚不赔'字样"}则终止推送。
# Dify 工作流中“话术合规检查”节点调用的 Flask 接口(简化版) from flask import Flask, request, jsonify import re app = Flask(__name__) PROHIBITED_WORDS = [ (r'稳赚不赔', '违反《金融营销宣传办法》第八条'), (r'保本|本金保障', '违反资管新规第十五条'), (r'预期收益.*?%', '需标注“业绩比较基准”,不可称“预期收益”') ] @app.route('/check', methods=['POST']) def check_text(): text = request.json.get('text', '') errors = [] for pattern, reason in PROHIBITED_WORDS: if re.search(pattern, text): errors.append(f"违规词:{pattern} —— {reason}") return jsonify({"pass": len(errors) == 0, "errors": errors})这个接口部署在内网,Dify 工作流通过 HTTP 节点调用。关键参数:Timeout=3s(防阻塞)、Retry=1(网络抖动重试)、Status Code=200 only(非 200 直接报错中断)。我们曾因超时设为 10s,导致单次策略推送平均耗时 12 秒,微信端显示“正在输入...”长达半分钟,用户流失率飙升 40%。
2.3 变量聚合器不是拼字符串,而是“策略上下文编织器”:解决多源数据时间戳错位问题
金融策略最头疼的是数据新鲜度不一致:用户风险测评是昨天做的,基金净值是今天早盘更新的,而持仓数据可能还卡在 T+1。Dify 变量聚合器(Variable Aggregator)在这里干一件关键事:给每个数据源打上可信时间戳,并按“最新可用”原则合成上下文。
例如,生成“股债平衡策略”时,需要:
- 用户风险等级(来源:CRM 系统,更新频率:实时)→ 时间戳
2025-04-05T10:23:15Z - 当前沪深300指数(来源:Wind API,更新频率:秒级)→
2025-04-05T10:23:18Z - 用户持有债券基金列表(来源:TA 系统,更新频率:T+1)→
2025-04-04T15:30:00Z
变量聚合器配置如下:
| 变量名 | 数据源 | 更新时间戳字段 | 容忍延迟 | 优先级 |
|---|---|---|---|---|
user_risk | CRM API | updated_at | 300s | 1 |
csi300 | Wind API | timestamp | 10s | 2 |
bond_holdings | TA API | trade_date | 86400s | 3 |
聚合逻辑:若bond_holdings时间戳距今 > 86400s(即超过 1 天),则该变量置空,工作流进入“数据陈旧”分支,返回“您的持仓数据暂未更新,建议稍后再试”。这避免了用昨日持仓计算今日策略的致命错误。这个配置在 Dify UI 里藏得深:进工作流 → 点击“变量聚合器”节点 → “高级设置” → 手动粘贴 JSON Schema 定义时间戳字段和容忍值。
3. MCP 协议不是玄学概念,而是解决“模型调用混乱”的手术刀:统一调度 Llama-3-70B、Qwen2-72B、本地风控模型三类引擎
很多团队卡在“为什么不用 Dify 直连模型?非要加一层 MCP?”——因为真实金融场景里,你绝不会只用一个模型。我们同时跑着:
- Llama-3-70B(阿里云百炼):生成自然语言策略话术,强在表达流畅;
- Qwen2-72B(魔搭 ModelScope):解析用户模糊需求,如“我想买点稳健的,别太折腾”,强在中文语义理解;
- 自研风控小模型(PyTorch,<100MB):轻量级二分类,判断“当前策略是否触发监管关注点”,如涉及杠杆、跨境、场外衍生品。
如果每个模型都单独写 API 调用,工作流里要堆 3 个 HTTP 节点,维护成本爆炸。MCP(Model Control Protocol)的价值,就是把它们抽象成统一的模型插座:Dify 只认 MCP 地址mcp://localhost:8000,背后由 MCP Server 动态路由到具体模型。这带来三个实打实的好处:
- 模型热替换:Qwen2-72B 出现 hallucination,运维一键切换到 Qwen2-57B,Dify 工作流完全无感;
- 调用审计:MCP Server 自动记录每次请求的
model_id,input_tokens,output_tokens,latency_ms,user_id,满足《证券期货业网络信息安全管理办法》第 32 条日志留存要求; - 降本增效:Llama-3-70B 按 token 计费贵,MCP Server 可配置缓存策略——相同用户+相同持仓数据,30 分钟内重复请求直接返回缓存结果,实测降低 37% 大模型调用成本。
3.1 本地部署 MCP Server:用 Python + FastAPI 三步搭起模型调度中枢
我们不用官方 MCP 参考实现(太重),而是用 120 行 Python 自建轻量 Server,核心就三个模块:
# mcp_server.py from fastapi import FastAPI, HTTPException from pydantic import BaseModel import httpx import asyncio from datetime import datetime app = FastAPI() # 模型路由表:model_id -> 实际 endpoint MODEL_ROUTES = { "llama3-70b": "https://dashscope.aliyuncs.com/api/v1/services/aigc/text-generation/generation", "qwen2-72b": "https://dashscope.aliyuncs.com/api/v1/services/aigc/text-generation/generation", "risk-guard": "http://localhost:8001/predict" # 本地风控模型 } class MCPRequest(BaseModel): model_id: str prompt: str parameters: dict = {} @app.post("/v1/chat/completions") async def mcp_chat(request: MCPRequest): if request.model_id not in MODEL_ROUTES: raise HTTPException(404, f"Model {request.model_id} not found") # 根据 model_id 路由到不同后端 if request.model_id in ["llama3-70b", "qwen2-72b"]: # 转换为 DashScope 格式 dashscope_payload = { "model": request.model_id, "input": {"messages": [{"role": "user", "content": request.prompt}]}, "parameters": request.parameters } async with httpx.AsyncClient() as client: resp = await client.post( MODEL_ROUTES[request.model_id], json=dashscope_payload, headers={"Authorization": "Bearer YOUR_DASHSCOPE_KEY"}, timeout=60.0 ) return resp.json() elif request.model_id == "risk-guard": # 直连本地风控模型(Flask) async with httpx.AsyncClient() as client: resp = await client.post( MODEL_ROUTES["risk-guard"], json={"text": request.prompt}, timeout=5.0 ) return {"choices": [{"message": {"content": resp.json()["result"]}}]}部署命令:
pip install fastapi uvicorn httpx uvicorn mcp_server:app --host 0.0.0.0 --port 8000 --workers 4注意:Dify 配置 MCP 地址时,必须用
http://开头,不能用mcp://(Dify 当前版本只支持 HTTP 协议的 MCP Server)。我们曾在此翻车:填mcp://localhost:8000导致 Dify 报错Invalid URL scheme,查日志才发现文档没写清楚。
3.2 Dify 中配置 MCP 模型:不是选“大模型”,而是填“MCP 地址”
在 Dify 控制台 → 设置 → 模型提供方 → 添加模型 → 选择 “MCP” 类型:
- 模型名称:
llama3-70b-finance(自定义,用于工作流中识别) - MCP 服务器地址:
http://host.docker.internal:8000(⚠️ 关键!Dify 运行在 Docker,localhost指向容器自身,必须用host.docker.internal指向宿主机) - 模型 ID:
llama3-70b(必须与 MCP Server 中MODEL_ROUTES键名一致) - 最大 Token:
4096 - 温度(Temperature):
0.3(金融话术需稳定,避免随机性)
配置完后,在工作流中拖入“大模型”节点,模型下拉框就会出现llama3-70b-finance。测试时,Dify 日志会打印Calling MCP server at http://host.docker.internal:8000/v1/chat/completions,看到这行就说明通了。
3.3 MCP 的真实价值:用“模型熔断”兜住突发流量——当 Qwen2-72B 崩溃时,自动切到备用模型
MCP Server 不只是路由,更是保险丝。我们在/v1/chat/completions接口里加了熔断逻辑:
# 在 mcp_server.py 中添加 from circuitbreaker import circuit @circuit(failure_threshold=5, recovery_timeout=60) async def call_qwen_model(payload): async with httpx.AsyncClient() as client: resp = await client.post( "https://dashscope.aliyuncs.com/api/v1/services/aigc/text-generation/generation", json=payload, headers={"Authorization": "Bearer KEY"}, timeout=30.0 ) if resp.status_code != 200: raise Exception(f"Qwen API error: {resp.status_code}") return resp.json() # 在主路由中调用 if request.model_id == "qwen2-72b": try: result = await call_qwen_model(dashscope_payload) except Exception as e: # 熔断触发,降级到 llama3-70b result = await call_llama_model(dashscope_payload) log_warning(f"Qwen2-72b failed, fallback to llama3-70b: {e}")效果:当 Qwen2-72B 因上游限流连续失败 5 次,MCP Server 自动将后续 60 秒内的请求全部路由到 Llama-3-70B,且记录告警日志。这避免了因单个模型故障导致整个策略推送服务雪崩。某次阿里云 DashScope 临时维护,我们靠此机制零感知扛过 47 分钟,用户无一投诉。
4. 微信端不是简单发消息,而是构建“可审计、可回溯、可灰度”的推送管道:从公众号模板消息到企微机器人全链路
“微信端自动化推送”常被误解为“用 itchat 发条消息”。但在金融合规框架下,这是一条必须留痕、可追溯、可开关的生产级管道。我们弃用所有个人号/模拟登录方案(不稳定、易封号、无审计),采用企业微信机器人 + 公众号模板消息双通道,并用 Kafka 做消息总线保证顺序与幂等。
4.1 为什么不用服务号订阅消息?——微信 2024 年新规下的硬约束
很多团队想用服务号“订阅消息”实现策略推送,但 2024 年 3 月微信发布《关于进一步规范金融类服务号模板消息使用的公告》,明确:
- 禁止使用“投资建议”“资产配置”“组合推荐”等词汇;
- 单用户每月最多接收 1 条非交易类通知;
- 必须提供“退订入口”,且退订后不得再次推送。
这意味着“每日策略推送”在服务号上直接被判死刑。我们转向两个合规路径:
- 企业微信:面向内部理财经理,推送“待跟进客户策略摘要”,属于 B2E 场景,无频次限制;
- 公众号模板消息:仅用于“用户主动触发后的单次反馈”,如用户发送“查看我的策略”,才回复一条带跳转链接的模板消息(链接指向 H5 策略详情页)。
4.2 企业微信机器人:用 Markdown 渲染策略卡片,支持一键转发给客户
企业微信机器人是我们的主力推送通道。关键不是发文字,而是用Markdown 卡片呈现专业感:
# 构建企业微信 Markdown 消息体(Python) def build_strategy_card(user_name, strategy_name, allocation, risk_note): return { "msgtype": "markdown", "markdown": { "content": f"""## 📈 {user_name} 的 {strategy_name} 策略({datetime.now().strftime('%m-%d')}) ✅ **配置建议** - 股票型基金:{allocation['equity']}% (推荐:华夏沪深300ETF联接A,代码 000051) - 债券型基金:{allocation['bond']}% (推荐:易方达稳健收益债券A,代码 110007) - 货币基金:{allocation['cash']}% ⚠️ **重要提示** {risk_note} 📎 [点击查看完整策略报告](https://your-domain.com/report?id=xxx) 💬 理财经理可直接转发此卡片给客户""" } } # 发送至企微机器人 webhook import requests webhook_url = "https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=YOUR_KEY" requests.post(webhook_url, json=build_strategy_card("张三", "股债平衡", {"equity":60,"bond":35,"cash":5}, "本策略历史年化波动率约 8.2%,适合风险承受能力为 C3 及以上的投资者"))效果:理财经理在企微收到一张带图标、分段、超链接的卡片,点击“转发给客户”即可原样发给微信好友,客户在微信里看到的就是富文本,体验接近 App。且所有推送记录存在企微后台,满足“操作留痕”监管要求。
4.3 公众号模板消息:用“一次一签”机制规避微信风控
用户在公众号内发送“策略”二字,触发策略生成。此时不能直接回复文字(会被限流),必须用模板消息。难点在于:微信模板消息必须提前在后台申请模板,且每个模板有唯一template_id,无法动态生成。
我们的解法是:预申请 10 个通用模板,用“动态占位符”覆盖所有场景。例如申请模板:
【${keyword1.DATA}】
${keyword2.DATA}
${keyword3.DATA}
${keyword4.DATA}
${keyword5.DATA}
然后在代码中动态填充:
# 公众号模板消息 payload template_data = { "first": {"value": "您的智能策略已生成", "color": "#173177"}, "keyword1": {"value": "股债平衡策略", "color": "#000000"}, "keyword2": {"value": "股票60% / 债券35% / 货币5%", "color": "#000000"}, "keyword3": {"value": "推荐产品:华夏沪深300ETF联接A(000051)", "color": "#000000"}, "keyword4": {"value": "风险提示:本策略历史年化波动率约 8.2%", "color": "#FF0000"}, "remark": {"value": "点击查看详情 →", "color": "#173177"} }这样,一个模板 ID 就能复用所有策略类型。我们预申请了 5 个模板(策略生成、调仓提醒、风险预警、产品到期、市场周报),覆盖 98% 场景。注意:template_id必须在微信公众号后台“功能->模板消息”里手动申请,Dify 无法代劳。
5. 避坑指南:Dify + MCP + 微信链路中,我们踩过的 4 个血泪坑与后悔药
这条链路看着是“Dify 工作流 → MCP Server → 微信 API”,实际落地时,80% 的故障发生在边界处。以下是我们在三家金融机构交付中反复验证的 4 个致命坑,附带可直接抄的解决方案。
5.1 坑:Dify 工作流中 HTTP 节点调用 MCP Server 返回 502,日志只显示“Connection refused”
- 现象:Dify 工作流执行到调用 MCP Server 的 HTTP 节点时失败,Dify 后台日志只有
HTTP request failed: Connection refused,MCP Server 日志完全空白。 - 原因:Dify 容器与 MCP Server 不在同一 Docker 网络,且
host.docker.internal在 Linux 宿主机上默认不可用(仅 macOS/Windows Docker Desktop 支持)。Dify 容器内ping host.docker.internal通不过。 - 解决:在
docker-compose.yml中显式声明网络别名:
然后 Dify 中 MCP 地址改为services: dify: # ... 其他配置 extra_hosts: - "host.docker.internal:host-gateway" # Linux 下强制映射 mcp-server: # ... 其他配置 networks: - dify-net networks: dify-net: driver: bridgehttp://mcp-server:8000(用服务名直连),彻底绕过host.docker.internal兼容性问题。
5.2 坑:微信模板消息发送成功,但用户收不到,后台显示“送达率 0%”
- 现象:调用微信模板消息 API 返回
{"errcode":0,"errmsg":"ok"},但用户手机无任何通知,公众号后台统计送达率为 0。 - 原因:微信要求模板消息必须在用户最近 7 天内主动与公众号互动过(发送消息、点击菜单、浏览图文),否则视为“无效用户”,不予推送。我们曾用测试号批量发给 1000 个老用户,其中 82% 超过 7 天未互动,全部被微信静默丢弃。
- 解决:在发送前,先调用微信 API 查询用户最后互动时间:
# 查询用户最近互动时间(需 access_token) url = f"https://api.weixin.qq.com/cgi-bin/user/info?access_token={token}&openid={openid}&lang=zh_CN" resp = requests.get(url).json() last_time = datetime.fromtimestamp(resp.get("last_time", 0)) if (datetime.now() - last_time).days > 7: # 改走企业微信通道,或发服务通知(需用户授权) send_to_work_wechat(openid, content)
5.3 坑:MCP Server 调用 DashScope API 时频繁超时,Dify 工作流卡在“运行中”状态
- 现象:Dify 工作流长时间显示“运行中”,MCP Server 日志显示
Read timeout on endpoint,DashScope 控制台无调用记录。 - 原因:DashScope 对免费额度用户限流极严,单 IP 每分钟最多 5 次请求。我们未配置重试退避,连续请求触发限流,后续请求全部超时。
- 解决:在 MCP Server 中加入指数退避重试:
import asyncio from tenacity import retry, stop_after_attempt, wait_exponential @retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=1, max=10)) async def call_dashscope_api(payload): async with httpx.AsyncClient() as client: resp = await client.post( "https://dashscope.aliyuncs.com/api/v1/services/aigc/text-generation/generation", json=payload, headers={"Authorization": "Bearer KEY"}, timeout=30.0 ) if resp.status_code == 429: # 限流 raise Exception("Rate limited") return resp.json()
5.4 坑:Dify 知识库更新后,工作流仍返回旧答案,疑似缓存未刷新
- 现象:修改了产品白名单 CSV,重新上传知识库并点击“重新索引”,但工作流调用时仍返回旧基金代码。
- 原因:Dify 知识库索引是异步任务,且默认缓存 1 小时。重新索引后,Dify 后台任务队列可能堆积,或缓存未失效。
- 解决:两步清空:
- 进入 Dify 后台 → 知识库 → 选择该知识库 → 点击“更多” → “清除缓存”;
- 在 Dify 容器内执行:
# 进入 Dify 容器 docker exec -it dify-web bash # 清除 Redis 缓存(Dify 默认用 Redis 存向量缓存) redis-cli -h redis -p 6379 FLUSHDB
6. 进阶技巧:用“策略灰度发布”代替“全量上线”——在 Dify 工作流里埋一个开关,让合规部门随时叫停
真正的金融系统上线,从来不是“一键发布”,而是“分批放量+实时监控+秒级熔断”。我们把这套机制直接嵌进 Dify 工作流,让合规同事不用懂技术,也能在后台点几下就控制策略推送范围。
6.1 灰度开关:用 Dify 环境变量 + 条件节点,实现“按用户 ID 尾号分流”
Dify 支持环境变量(Environment Variables),我们创建一个名为STRATEGY_GRAYSCALE_RATE的变量,值为0.1(表示 10% 用户灰度)。在工作流中插入一个“条件节点”:
- 条件表达式:
{{ user_id | last_digit }} % 10 < {{ env.STRATEGY_GRAYSCALE_RATE * 10 | int }} - True 分支:走完整策略生成 → 合规检查 → 微信推送;
- False 分支:返回固定话术“策略服务升级中,敬请期待”。
其中last_digit是我们自定义的 Jinja2 过滤器(需在 Dify 后端代码中注册):
# 在 Dify 源码的 jinja2_env.py 中添加 def last_digit(value): try: return int(str(value)[-1]) except: return 0 env.filters['last_digit'] = last_digit这样,当STRATEGY_GRAYSCALE_RATE=0.3时,用户 ID 尾号为 0/1/2 的用户进入灰度;设为1.0则全量。合规同事只需登录 Dify 后台 → 设置 → 环境变量 → 修改数值 → 保存,无需重启服务。
6.2 实时监控看板:用 Kafka + Grafana 拉出“策略推送健康度仪表盘”
我们把每一次策略推送的关键事件发到 Kafka:
strategy.request:用户触发,含user_id,trigger_type,timestampstrategy.generated:策略生成成功,含allocation,model_used,latency_msstrategy.sent:微信推送成功,含channel(wechat/qrwork),status
用 Grafana 连接 Kafka,建一个看板,核心指标:
| 指标 | 查询语句(Prometheus + Kafka Exporter) | 业务意义 |
|---|---|---|
| 推送成功率 | rate(kafka_topic_partition_current_offset{topic="strategy.sent"}[5m]) / rate(kafka_topic_partition_current_offset{topic="strategy.request"}[5m]) | 低于 99.5% 触发告警 |
| 模型平均延迟 | histogram_quantile(0.95, sum(rate(kafka_consumer_fetch_latency_seconds_bucket{topic="strategy.generated"}[5m])) by (le)) | 超过 8s 需扩容 MCP Server |
| 灰度覆盖率 | count by (gray_flag) (kafka_topic_partition_current_offset{topic="strategy.request"}) | 确认灰度比例是否符合配置 |
当某天下午 2 点推送成功率跌到 92%,看板立刻标红,运维查 Kafka 发现是 DashScope 限流,立即把灰度率从 1.0 降到 0.3,10 秒内恢复。这种“可观测性”不是锦上添花,而是金融系统的生命线。
6.3 合规回溯:用 Dify 的“会话溯源”功能,一键导出某次推送的完整决策链
监管检查时,最常问:“这个策略是怎么生成的?依据哪些数据?谁审核的?” Dify 天然支持会话溯源。我们做了两件事:
- 在工作流每个关键节点加日志记录:用“HTTP 节点”调用内部日志服务,传入
session_id,node_name,input,output,timestamp; - 导出为 PDF 报告:Dify API 提供
/sessions/{id}/messages接口,我们写了个脚本,输入用户 ID 和日期,自动拉取该用户所有策略会话,渲染成带时间轴、节点截图、原始数据的 PDF,加盖电子章。
# 生成合规报告的伪代码 def generate_compliance_report(user_id, date): session_ids = get_sessions_by_user(user_id, date) # 从 Dify API 拉 report_data = [] for sid in session_ids: messages = dify_api.get_session_messages(sid) for msg in messages: if msg.role == "assistant" and "策略" in msg.content: # 解析策略 JSON 片段 strategy_json = extract_strategy_json(msg.content) report_data.append({ "timestamp": msg.created_at, "strategy": strategy_json, "data_sources": get_data_sources(strategy_json), # 从变量聚合器日志反查 "model_used": get_model_from_log(msg.created_at) }) return render_pdf(report_data)这份报告,就是我们应对现场检查的“后悔药”。去年某次证监局抽查,我们 3 分钟内提供了张三客户 3 月 15 日策略推送的完整决策链,从用户提问、风险等级查询、模型调用、合规检查到微信发送,全部可追溯,当场过关。
我带过的每个金融项目,上线前都坚持做一件事:把 Dify 工作流截图、MCP Server 日志片段、微信推送回执、Kafka 监控截图,四张图拼成一页 A4,贴在工位墙上。不是为了好看,是提醒自己——我们写的不是代码,是客户账户里的真金白银,是监管罚单上的数字,是理财经理每天面对客户的底气。希望帮到你。
本文还有配套的精品资源,点击获取