企业Agent实战搭建全流程3/3——药企多Agent跨部门协作

这是"药企AI Agent落地三部曲"的第3篇,也是最后一篇。
-
第1篇:文献调研与竞品情报——信息淹没,人工整理效率极低
-
第2篇:临床方案与注册文件——长文档修改联动,格式校验耗时
-
本篇:跨部门数据孤岛——多部门数据不通,协作靠手工转格式
本三部曲完整记录了从单Agent逐步扩展到多Agent协作系统的过程。所有工具全线开源,可私有化部署。
本篇在第1、2篇的架构基础上,展示了用Dify原生的多Agent能力,把系统从"单Agent效率工具"升级为"跨部门协作系统"。新增的脱敏和审计组件全部开源,可私有化部署。
1. 问题描述
单Agent在各业务线跑起来之后,部门的效率确实提升了。但是药企的业务链条是跨部门的——临床运营、注册申报、医学事务,三个部门的数据是串在一起的,不是各干各的。
举几个实际遇到的情况:
-
临床运营的数据在CTMS里,注册申报需要eTMF格式,医学报告又是一套格式。同一批临床试验数据,三个部门用三种格式存着,每次跨部门对接要花3天,其中2天是花在格式转换上
-
每月平均发现15-20处跨部门数据矛盾。比如临床运营那边更新了一个站点的入组人数,注册部的文件里还是老数字,两边对不上
-
临床方案改了一个入组标准,注册部3天后才知道。因为方案变更的通知走的是邮件+Excel,没有人跟踪到底谁看了谁没看
-
一个III期试验入组完成,要在CTMS、eTMF、医学报告三个系统里分别手动更新数据。每次更新都怕出错,出错之后又要花时间去对
再说说硬约束,这些是药企特有的,也是跨部门协作最头疼的地方:
-
患者信息绝对不能跨部门明文流转。这是GxP的红线,碰了就是重大偏差。医学部可以看到临床数据的汇总结果,但是不能看到患者的明文信息;注册部需要引用临床数据写申报文件,但是受试者的个人信息必须脱敏
-
审计日志必须完整可追溯。GxP审查的硬指标,每一次数据访问、每一次数据修改都要有记录。出了问题能倒查是谁在什么时间做了什么操作
-
各部门数据访问权限必须隔离。临床运营可以看到所有站点的所有数据,医学部只能看到脱敏后的汇总数据,注册部只能看到跟申报相关的内容。权限不能混
-
复用前面已经部署的基础设施。Dify、Qwen3.5、PGVector、Mem0这些都已经在跑了,不能推翻重来
说白了,你要在"患者信息不泄露 + 审计日志全覆盖 + 部门权限隔离 + 复用已有架构"这四个约束条件下,让数据在部门之间安全地流动起来。
注意,这里说的是"安全地流动",不是"随便流动"。药企的数据治理跟互联网公司完全不一样,不能因为上了AI系统就把合规底线给突破了。
2. 方案设计
核心思路:从"单Agent效率工具"升级为"多Agent协作系统"。
这里有一个关键的判断要先说清楚:多Agent系统的目标不是"打破部门墙"。部门墙是组织问题,技术解决不了。你要做的是"让Agent适配部门边界"——每个Agent在自己的部门范围内做擅长的事,通过安全的方式交换必要的信息。
设计原则就八个字:数据可用不可见,操作可溯可审计。
翻译成技术语言就是:跨部门传递的数据必须经过脱敏,所有的访问和修改操作必须记入审计日志。
因为第1、2篇的基础架构已经在跑了(Dify、Qwen3.5、PGVector、RustFS、BGE-Reranker、Mem0),所以这篇只说新增的部分。
新增技术选型
|
新增组件 |
选型 |
为什么选它 |
|---|---|---|
|
多Agent编排 |
Dify Workflow路由 |
复用Dify,原生支持多Agent条件分支和路由,不需要额外引入新平台 |
|
共享记忆层 |
Mem0扩展(摘要层+详情层) |
复用已有Mem0,扩展实现跨Agent上下文传递,不需要再部署一套新的记忆系统 |
|
数据脱敏 |
Qwen-NER |
复用已经部署的Qwen3.5做命名实体识别,零额外部署成本 |
|
审计日志 |
EasySearch |
INFINI Labs开源的日志分析引擎,兼容Elasticsearch API,本地化部署运维可控,GxP审计友好 |
|
权限控制 |
Nginx RBAC + Dify API权限 |
各部门数据访问权限隔离,Nginx做网关层,Dify做应用层,双重保障 |
|
格式映射 |
pandas + 自建映射规则 |
跨部门数据格式自动转换,CTMS→eTMF→医学报告的字段映射 |
为什么用EasySearch不直接用Elasticsearch?因为EasySearch是INFINI Labs开发的国产开源搜索引擎,完全兼容Elasticsearch的API,但是部署更轻量、运维更简单。Docker一行命令就能跑起来,而且社区文档对中文场景支持更好。审计日志这种场景不需要Elasticsearch的全部功能,EasySearch完全够用。
Qwen-NER这里要说明一下:它不是一个单独的模型,而是用已经部署的Qwen3.5通过特定的prompt来做命名实体识别任务。后面Step 3会详细说。如果后期对NER精度要求更高,可以用少量标注数据微调一个专用的NER模型,但是初期用prompt方式就够了。
架构图

跟第1篇的架构图对比一下,变化其实不大——LLM底座没变,向量数据库没变,记忆层在原有的Mem0基础上做了扩展,然后新增了脱敏、权限网关、审计日志三个组件。整个架构是一步步长出来的,不是第一天就规划成这样的。
3. 实施过程
下面按天拆解。第1篇已经讲过的Dify部署、Qwen3.5部署、PGVector配置这些基础部分不再重复,这篇只说新增和变化的部分。
Step 1:多Agent职责定义与创建(Day 1-3)
这一步要做的事情是:在Dify中创建3个独立的Agent,每个Agent负责一个业务线,职责边界必须清晰。
为什么要拆成3个Agent而不是用1个Agent做所有事?一开始确实是这么想的——搞一个"全能Agent",把所有部门的事情都交给它。结果发现不行,prompt写了一大堆,各种条件和权限判断嵌套在一起,Agent自己都搞不清楚该用哪个工具。而且一个Agent的上下文窗口装不下三个部门的知识库和工具集。
所以最后决定拆开。每个Agent有明确的职责范围、独立的system prompt、专属的工具集和知识库、独立的权限范围。
三个Agent的定义:
医学部Agent——文献检索、竞品情报、医学报告生成。能力范围是公开文献和内部医学知识库,不能访问患者数据。这个Agent其实就是第1篇搭建的那个,直接复用。
注册部Agent——注册文件生成、格式校验、合规检查。在第2篇的长文档处理Agent基础上扩展,增加了eTMF格式转换和合规校验工具。temperature设的是0.1,因为注册文件要求精确,不能有"创造性"输出。
临床运营Agent——临床试验数据管理、跨部门协调、方案变更通知。这个Agent能访问CTMS数据,但是输出给其他部门的数据必须经过脱敏。
下面是通过Dify API创建这三个Agent的代码:
import requests
import json
DIFY_BASE_URL = "http://your-dify-server/v1"
DIFY_API_KEY = "your-dify-admin-api-key"
headers = {
"Authorization": f"Bearer {DIFY_API_KEY}",
"Content-Type": "application/json"
}
# 三个Agent的配置,每个都有独立的system prompt和工具集
agents_config = {
"medical_affairs": {
"name": "医学部Agent",
"description": "负责文献检索、竞品情报、医学报告生成",
"system_prompt": """你是药企医学部的AI助手。你的职责包括:
1. 文献检索与情报整理:检索PubMed、ClinicalTrials.gov等数据库,整理竞品动态
2. 医学报告生成:根据临床数据生成医学报告,使用标准医学术语
3. 知识管理:维护医学知识库,确保信息的准确性和时效性
你的权限范围:
- 可以访问:公开文献数据库、内部医学知识库(PGVector中medical_kb命名空间)
- 不能访问:患者个人信息、受试者编号、临床运营原始数据
- 输出要求:涉及临床数据时只使用汇总统计,不引用个案信息""",
"temperature": 0.7,
"tools": ["pubmed_search", "knowledge_base_qa", "report_generator"],
"knowledge_bases": ["medical_literature", "competitor_intelligence"]
},
"regulatory": {
"name": "注册部Agent",
"description": "负责注册文件生成、格式校验、合规检查",
"system_prompt": """你是药企注册部的AI助手。你的职责包括:
1. 注册文件生成:根据临床数据生成CTD格式的注册申报文件
2. 格式校验:检查文件格式是否符合eTMF/CTD规范要求
3. 合规检查:验证文件内容是否满足ICH指南和NMPA要求
你的权限范围:
- 可以访问:注册知识库、已脱敏的临床数据汇总、eTMF模板库
- 不能访问:患者个人信息、未脱敏的临床数据
- 输出要求:所有输出必须精确、可追溯,引用数据必须标注来源和版本号
- 特别注意:注册文件不允许有任何推测性或创造性内容,所有表述必须有数据支撑""",
"temperature": 0.1,
"tools": ["ctd_generator", "format_validator", "compliance_checker"],
"knowledge_bases": ["regulatory_guidelines", "etmf_templates"]
},
"clinical_ops": {
"name": "临床运营Agent",
"description": "负责临床试验数据管理、跨部门协调、方案变更通知",
"system_prompt": """你是药企临床运营部的AI助手。你的职责包括:
1. 数据管理:管理CTMS中的临床试验数据,包括入组进度、站点状态、数据质量
2. 跨部门协调:当数据变更影响到其他部门时,触发通知和协作流程
3. 方案变更管理:跟踪临床方案的变更记录,评估影响范围,通知相关部门
你的权限范围:
- 可以访问:CTMS全部数据、站点数据、入组数据
- 不能访问:医学部内部知识库、注册部内部文档
- 输出要求:向其他部门传递数据时,必须先通过脱敏检查
- 受试者姓名、身份证号、联系方式:必须删除
- 受试者编号(S001-S999格式):必须替换为随机编码
- 具体入组日期:仅保留月份,不保留具体日期
- 特别注意:涉及患者信息的操作必须记录审计日志""",
"temperature": 0.3,
"tools": ["ctms_query", "data_export", "change_notifier", "desensitization_check"],
"knowledge_bases": ["clinical_protocols", "site_management"]
}
}
# 通过Dify API逐个创建Agent
created_agents = {}
for agent_key, config in agents_config.items():
payload = {
"name": config["name"],
"description": config["description"],
"mode": "agent-chat",
"model_config": {
"provider": "tongyi", # 对接已部署的Qwen3.5
"model": "qwen3.5-35b-a3b",
"temperature": config["temperature"],
"completion_params": {
"max_tokens": 4096
}
},
"pre_prompt": config["system_prompt"],
"user_input_form": [],
"dataset_ids": config.get("knowledge_bases", []),
"agent_config": {
"tools": config["tools"]
}
}
response = requests.post(
f"{DIFY_BASE_URL}/apps",
headers=headers,
json=payload
)
if response.status_code == 200:
agent_data = response.json()
created_agents[agent_key] = agent_data["id"]
print(f"创建成功:{config['name']},ID:{agent_data['id']}")
else:
print(f"创建失败:{config['name']},错误:{response.text}")
print(f"\n已创建的Agent映射:{json.dumps(created_agents, ensure_ascii=False, indent=2)}")
创建完Agent之后,还有一个关键步骤——给每个Agent配置工具。以临床运营Agent的脱敏检查工具为例:
# 在Dify中注册自定义工具(通过Dify的自定义工具API)
# 这个工具会被临床运营Agent调用,用于在数据输出前做脱敏检查
desensitization_tool_config = {
"name": "desensitization_check",
"description": "在数据输出前进行脱敏检查,识别并替换患者个人信息",
"parameters": {
"type": "object",
"properties": {
"text": {
"type": "string",
"description": "需要脱敏检查的文本内容"
},
"target_department": {
"type": "string",
"enum": ["medical", "regulatory", "internal"],
"description": "数据接收部门,不同部门脱敏级别不同"
}
},
"required": ["text", "target_department"]
}
}
print("脱敏检查工具配置已准备,需要在Dify自定义工具页面完成注册")
⚠️ 踩坑记录
坑1:一开始想让一个Agent做所有事,把三个部门的prompt都塞进同一个Agent里。结果发现prompt超过3000字之后,Agent的指令遵循能力明显下降——该调工具的时候不调,不该输出的时候乱输出。而且不同部门的权限要求是冲突的(一个部门能看的东西另一个部门不能看),在同一个Agent里根本处理不了。最后老老实实拆成了3个。
坑2:注册部Agent的temperature一开始默认设了0.7,结果生成的注册文件里出现了一些"看起来合理但是找不到数据来源"的表述。这在注册场景里是致命的——每一句话都必须有数据出处。把temperature调到0.1之后,输出变得非常保守和精确,再也没出现过"创造性发挥"的问题。
📎 官方文档:[Dify Agent配置](https://docs.dify.ai/guides/application-orchestrate/agent)
Step 2:共享记忆层搭建(Day 4-7)
三个Agent各自有了,但是它们之间需要交换信息。比如临床运营Agent发现入组进度延迟了,这个信息需要让注册部Agent知道(可能需要调整申报时间线),也需要让医学部Agent知道(可能影响医学报告的数据截止点)。
但是不能把所有信息都共享——患者信息、站点内部数据这些是不能跨部门看到的。所以需要一个"摘要层+详情层"的双层记忆架构。
第1篇已经部署了Mem0,这一步是在Mem0基础上做扩展。
设计思路:
-
摘要层:经过脱敏的事件摘要,所有Agent都能读。比如"XX试验入组进度延迟约2周,原因是XX站点启动延迟"。这里面没有患者信息,没有受试者编号
-
详情层:完整数据,只有对应部门的Agent能读写。比如"XX试验S003号受试者因为AE退出",这种信息只有临床运营Agent能写入和读取

from mem0 import Memory
from openai import OpenAI
import json
import re
from typing import Optional
# 复用已部署的Qwen3.5作为LLM和Embedding
client = OpenAI(
base_url="http://your-qwen-server:8000/v1",
api_key="not-needed" # 私有化部署不需要真实API Key
)
# ---- 摘要层配置 ----
# 所有Agent共享的脱敏记忆,存放事件摘要
summary_memory = Memory.from_config({
"vector_store": {
"provider": "pgvector",
"config": {
"collection_name": "shared_summary_memory",
"host": "your-pg-host",
"port": 5432,
"dbname": "agent_memory",
"user": "agent_user",
"password": "your-password"
}
},
"llm": {
"provider": "openai",
"config": {
"model": "qwen3.5-35b-a3b",
"temperature": 0.1,
"max_tokens": 2048,
"api_key": "not-needed",
"openai_base_url": "http://your-qwen-server:8000/v1"
}
},
"embedder": {
"provider": "openai",
"config": {
"model": "qwen3.5-35b-a3b",
"api_key": "not-needed",
"openai_base_url": "http://your-qwen-server:8000/v1"
}
},
# 摘要层的相关度阈值设低一些(0.3),因为跨部门查询的表述差异比同部门大
# 比如临床运营说"入组延迟",注册部可能搜"申报时间线调整"
"retrieval_config": {
"top_k": 5,
"score_threshold": 0.3
}
})
# ---- 详情层配置 ----
# 每个部门独立的记忆空间,只有本部门的Agent能访问
# 以临床运营为例,其他部门类似,只是collection_name和agent_id不同
detail_memory_clinical = Memory.from_config({
"vector_store": {
"provider": "pgvector",
"config": {
"collection_name": "detail_memory_clinical_ops",
"host": "your-pg-host",
"port": 5432,
"dbname": "agent_memory",
"user": "agent_user",
"password": "your-password"
}
},
"llm": {
"provider": "openai",
"config": {
"model": "qwen3.5-35b-a3b",
"temperature": 0.1,
"max_tokens": 2048,
"api_key": "not-needed",
"openai_base_url": "http://your-qwen-server:8000/v1"
}
},
"embedder": {
"provider": "openai",
"config": {
"model": "qwen3.5-35b-a3b",
"api_key": "not-needed",
"openai_base_url": "http://your-qwen-server:8000/v1"
}
},
"retrieval_config": {
"top_k": 10,
"score_threshold": 0.5
}
})
# 注册部和医学部的详情层配置结构一样,只改collection_name
# detail_memory_regulatory = Memory.from_config({...collection_name: "detail_memory_regulatory"...})
# detail_memory_medical = Memory.from_config({...collection_name: "detail_memory_medical"...})
# ---- 核心功能函数 ----
def write_to_shared_summary(agent_id: str, event_text: str):
"""
将事件写入共享摘要层。
写入前自动脱敏,确保不包含患者个人信息。
"""
# 第一步:用Qwen做脱敏检查
desensitized = desensitize_text(event_text)
# 第二步:生成事件摘要
summary_response = client.chat.completions.create(
model="qwen3.5-35b-a3b",
messages=[{
"role": "system",
"content": "你是一个事件摘要生成器。将以下事件描述压缩为一句话摘要,只保留关键信息(时间、事件、影响),删除所有个人信息和可追溯标识。"
}, {
"role": "user",
"content": desensitized
}],
temperature=0.1
)
summary = summary_response.choices[0].message.content
# 第三步:写入摘要层记忆
summary_memory.add(
data=summary,
user_id=f"shared_event",
metadata={
"source_agent": agent_id,
"event_type": "cross_department",
"original_length": len(event_text)
}
)
print(f"[摘要层写入] Agent={agent_id}, 摘要:{summary[:80]}...")
return summary
def write_to_detail(agent_id: str, detail_text: str, department: str):
"""
将完整数据写入对应部门的详情层。
只有对应部门的Agent能读写。
"""
memory_map = {
"clinical_ops": detail_memory_clinical,
"regulatory": detail_memory_regulatory,
"medical": detail_memory_medical
}
target_memory = memory_map.get(department)
if not target_memory:
raise ValueError(f"未知部门:{department}")
target_memory.add(
data=detail_text,
user_id=agent_id,
metadata={
"department": department,
"access_level": "restricted",
"timestamp": "auto"
}
)
print(f"[详情层写入] Agent={agent_id}, 部门={department}")
def read_shared_summary(query: str, top_k: int = 5):
"""
从共享摘要层读取信息。所有Agent都可以调用。
"""
results = summary_memory.search(query=query, top_k=top_k)
return results
def read_detail(query: str, department: str, agent_id: str, top_k: int = 10):
"""
从详情层读取信息。只能读取本部门的详情。
调用时会校验agent_id是否有权限访问该部门的数据。
"""
# 权限校验:agent_id必须属于对应部门
permission_map = {
"clinical_ops_agent": "clinical_ops",
"regulatory_agent": "regulatory",
"medical_agent": "medical"
}
allowed_department = permission_map.get(agent_id)
if allowed_department != department:
raise PermissionError(
f"Agent {agent_id} 无权访问 {department} 部门的详情数据"
)
memory_map = {
"clinical_ops": detail_memory_clinical,
"regulatory": detail_memory_regulatory,
"medical": detail_memory_medical
}
target_memory = memory_map.get(department)
results = target_memory.search(query=query, top_k=top_k)
return results
# ---- 脱敏函数(简化版,完整版在Step 3) ----
def desensitize_text(text: str) -> str:
"""
对文本做基础脱敏。Step 3会用Qwen-NER替换这个简化版本。
"""
# 身份证号
text = re.sub(r'\d{17}[\dXx]', '[已脱敏]', text)
# 手机号
text = re.sub(r'1[3-9]\d{9}', '[已脱敏]', text)
# 受试者编号(S开头+3位数字)
text = re.sub(r'S\d{3}', '[受试者编号已脱敏]', text)
# 姓名模式(2-4个汉字,在"姓名:"后面)
text = re.sub(r'(姓名[::])\s*[\u4e00-\u9fff]{2,4}', r'\1[已脱敏]', text)
return text
# ---- 使用示例 ----
# 场景:临床运营Agent发现入组进度延迟,需要通知其他部门
# 1. 写入详情层(包含完整信息,只有临床运营能看)
write_to_detail(
agent_id="clinical_ops_agent",
detail_text="试验C-2024-003的S015号受试者,张三,身份证号110101199001011234,因不良事件于2024年3月15日退出试验。该站点本月入组目标差2例。",
department="clinical_ops"
)
# 2. 写入摘要层(脱敏后,所有部门都能看到)
write_to_shared_summary(
agent_id="clinical_ops_agent",
event_text="试验C-2024-003的S015号受试者,张三,身份证号110101199001011234,因不良事件于2024年3月15日退出试验。该站点本月入组目标差2例。"
)
# 摘要层实际存储的内容类似:
# "试验C-2024-003某站点因受试者退出导致本月入组进度延迟,缺口2例。"
# 3. 注册部Agent查询相关事件
results = read_shared_summary("C-2024-003 入组进度")
print(f"注册部查到的摘要:{results}")
# 4. 注册部Agent如果想看详情,会被拒绝(权限隔离)
try:
detail = read_detail("C-2024-003 受试者退出", "clinical_ops", "regulatory_agent")
except PermissionError as e:
print(f"权限拒绝:{e}")
# 输出:权限拒绝:Agent regulatory_agent 无权访问 clinical_ops 部门的详情数据
⚠️ 踩坑记录
坑1:一开始直接把全部记忆都共享给所有Agent,觉得这样信息传递效率最高。结果测试的时候发现,临床运营Agent的记忆里有受试者的真实姓名和身份证号,注册部Agent能直接搜到。这在GxP审查里是重大偏差——患者信息跨部门明文流转了。后来才改成双层架构,摘要层经过脱敏才能共享,详情层严格隔离。
坑2:摘要层的相关度阈值一开始用的是0.5(跟第1篇一样)。但是跨部门查询的时候,同一个事件的表述差异很大——临床运营说"入组延迟",注册部搜"申报时间线",医学部搜"数据截止"。0.5的阈值导致很多相关摘要搜不出来。调到0.3之后,召回率明显提升。但是调低阈值也带来了噪声,所以配合summarymemory的metadata过滤(按eventtype、source_agent筛选)来控制精度。
📎 官方文档:[Mem0](https://docs.mem0.ai/)
Step 3:数据脱敏与权限网关(Day 8-10)
Step 2里的脱敏函数是简化版,用的是正则匹配。这种方式能处理通用的PII(身份证号、手机号),但是处理不了药企专有的敏感信息,比如受试者编号(S001-S999格式)。而且正则的误报率比较高,比如把产品批号误识别为受试者编号。
这一步要用Qwen3.5做NER(命名实体识别),让模型来识别哪些是敏感信息。因为Qwen3.5已经在跑了,所以不需要额外部署任何东西,只需要写一个NER的prompt模板。

import json
import re
from openai import OpenAI
client = OpenAI(
base_url="http://your-qwen-server:8000/v1",
api_key="not-needed"
)
class QwenNERDesensitizer:
"""
用Qwen3.5做命名实体识别(NER),识别并替换文本中的敏感信息。
这里不是用一个单独的NER模型,而是通过prompt让Qwen3.5执行NER任务。
好处是零额外部署成本,坏处是速度比专用模型慢一些。
如果后期对速度有要求,可以用50-100条标注数据微调一个专用NER模型。
"""
def __init__(self):
# 定义药企场景的敏感实体类型
self.entity_types = [
"PATIENT_NAME", # 患者姓名
"ID_NUMBER", # 身份证号
"PHONE_NUMBER", # 手机号
"SUBJECT_ID", # 受试者编号(S001-S999格式)
"DATE_OF_BIRTH", # 出生日期
"ADDRESS", # 详细地址
"HOSPITAL_NAME", # 具体医院名称(某些场景下需要脱敏)
"EXACT_DATE" # 精确日期(在某些跨部门场景需要模糊化)
]
def identify_entities(self, text: str) -> list:
"""
用Qwen3.5识别文本中的敏感实体。
返回实体列表,每个实体包含类型、位置和原文。
"""
entity_type_list = "\n".join([f"- {et}" for et in self.entity_types])
prompt = f"""你是一个药企数据脱敏专家。请识别以下文本中的敏感实体。
敏感实体类型包括:
{entity_type_list}
注意以下药企特有的识别规则:
1. 受试者编号格式为S后面跟3位数字(如S001、S123、S999),包括小写s开头
2. 受试者编号也包含SCREEN-XXX格式(如SCREEN-001)
3. 临床试验编号格式为C-YYYY-NNN(如C-2024-003)不属于敏感信息,不需要脱敏
4. 药物编号(如D-001)不属于敏感信息
请以下面的JSON格式返回结果,如果没有找到某类实体则不要包含在结果中:
[{"type": "实体类型", "value": "原始文本", "start": 起始位置, "end": 结束位置}]
如果文本中没有敏感信息,返回空列表 []。
文本:
{text}"""
response = client.chat.completions.create(
model="qwen3.5-35b-a3b",
messages=[{"role": "user", "content": prompt}],
temperature=0.0, # NER任务需要确定性输出,temperature设为0
max_tokens=1024
)
result_text = response.choices[0].message.content.strip()
# 提取JSON部分(模型有时候会在JSON外面加一些说明文字)
json_match = re.search(r'\[.*\]', result_text, re.DOTALL)
if json_match:
try:
entities = json.loads(json_match.group())
return entities
except json.JSONDecodeError:
print(f"NER结果解析失败:{result_text[:200]}")
return []
return []
def desensitize(self, text: str, target_department: str = "shared") -> str:
"""
对文本进行脱敏处理。
根据目标部门决定脱敏级别。
"""
entities = self.identify_entities(text)
if not entities:
return text
# 按位置从后往前替换,避免位置偏移
entities_sorted = sorted(entities, key=lambda x: x["start"], reverse=True)
for entity in entities_sorted:
entity_type = entity["type"]
start = entity["start"]
end = entity["end"]
# 根据实体类型和目标部门决定替换策略
if entity_type in ["PATIENT_NAME", "ID_NUMBER", "PHONE_NUMBER", "ADDRESS"]:
# 这些在任何跨部门场景都要完全替换
replacement = "[已脱敏]"
elif entity_type == "SUBJECT_ID":
# 受试者编号:内部保留,跨部门替换
if target_department == "shared":
replacement = "[受试者编号已脱敏]"
else:
continue # 内部使用,不替换
elif entity_type == "DATE_OF_BIRTH":
replacement = "[出生日期已脱敏]"
elif entity_type == "EXACT_DATE":
if target_department == "shared":
# 跨部门只保留月份
original = entity["value"]
month_match = re.match(r'(\d{4}年\d{1,2}月)', original)
if month_match:
replacement = month_match.group(1)
else:
replacement = "[日期已模糊化]"
else:
continue
else:
replacement = "[已脱敏]"
text = text[:start] + replacement + text[end:]
return text
# ---- 使用示例 ----
desensitizer = QwenNERDesensitizer()
# 测试文本
test_text = """
受试者S015,张三,男,身份证号110101199001011234,
出生日期1990年1月1日,于2024年3月15日因不良事件退出试验。
联系电话13812345678,居住地址为北京市朝阳区XX路XX号。
该受试者属于试验C-2024-003,使用药物D-001。
"""
# 跨部门脱敏(脱敏级别最高)
desensitized_shared = desensitizer.desensitize(test_text, target_department="shared")
print("跨部门脱敏结果:")
print(desensitized_shared)
# 输出中:张三→[已脱敏],身份证号→[已脱敏],S015→[受试者编号已脱敏]
# 13812345678→[已脱敏],地址→[已脱敏],2024年3月15日→2024年3月
# 试验C-2024-003和药物D-001保留(非敏感信息)
# ---- Nginx权限网关配置 ----
#
# 下面是Nginx的RBAC配置,作为权限网关层。
# 不同部门的API请求,通过Nginx注入不同的权限头。
# Dify内部会读取这些权限头做二次校验。
#
# 配置文件:/etc/nginx/conf.d/agent_gateway.conf
下面是Nginx权限网关的配置文件:
# 生成Nginx配置文件的Python脚本
# 实际使用时把输出内容保存到 /etc/nginx/conf.d/agent_gateway.conf
nginx_config = """
# 药企多Agent系统 - Nginx权限网关配置
# 不同部门通过不同的URL路径访问,Nginx自动注入对应的权限头
# 医学部的API入口
server {
listen 8443 ssl;
server_name agent-gateway.pharma-internal.com;
ssl_certificate /etc/nginx/ssl/agent-gateway.crt;
ssl_certificate_key /etc/nginx/ssl/agent-gateway.key;
# 医学部 - 只能访问医学相关的Agent API
location /api/medical/ {
# 注入医学部的权限头
proxy_set_header X-Department "medical";
proxy_set_header X-Data-Scope "literature,intelligence,report";
proxy_set_header X-Access-Level "read_summary,read_own_detail";
# 转发到Dify API
proxy_pass http://127.0.0.1:3000/;
# 限速,防止大量请求冲击
limit_req zone=api_limit burst=20 nodelay;
}
# 注册部 - 只能访问注册相关的Agent API
location /api/regulatory/ {
proxy_set_header X-Department "regulatory";
proxy_set_header X-Data-Scope "filing,compliance,etmf";
proxy_set_header X-Access-Level "read_summary,read_own_detail";
proxy_pass http://127.0.0.1:3000/;
limit_req zone=api_limit burst=20 nodelay;
}
# 临床运营 - 可以访问全部Agent API,但是输出受脱敏约束
location /api/clinical/ {
proxy_set_header X-Department "clinical_ops";
proxy_set_header X-Data-Scope "ctms,site,enrollment,all_summary";
proxy_set_header X-Access-Level "read_summary,read_own_detail,write_summary";
proxy_pass http://127.0.0.1:3000/;
limit_req zone=api_limit burst=30 nodelay;
}
# 共享的只读摘要查询接口(所有部门都能访问)
location /api/shared/summary/ {
proxy_set_header X-Department "$http_x_department";
proxy_set_header X-Data-Scope "read_summary_only";
proxy_set_header X-Access-Level "read_summary";
proxy_pass http://127.0.0.1:3000/;
limit_req zone=api_limit burst=10 nodelay;
}
}
# 限速区域定义
limit_req_zone $binary_remote_addr zone=api_limit:10m rate=10r/s;
# 访问日志 - 记录所有请求,供审计使用
log_format audit_log '$remote_addr - $remote_user [$time_local] '
'"$request" $status $body_bytes_sent '
'"$http_x_department" "$http_x_data_scope" '
'"$http_x_access_level"';
access_log /var/log/nginx/agent_audit.log audit_log;
"""
# 保存配置文件
with open("/etc/nginx/conf.d/agent_gateway.conf", "w") as f:
f.write(nginx_config)
print("Nginx权限网关配置已生成")
print("重启Nginx使配置生效:nginx -s reload")
如果后期对NER精度要求更高(比如需要识别更多药企专有格式),可以用少量标注数据微调一个专用的NER模型。下面是微调数据准备和训练的代码:
# ---- 可选:用标注数据微调NER专用模型 ----
# 如果prompt方式的NER识别率不够(比如受试者编号识别不全),
# 可以用50-100条标注数据微调一个专用模型。
# 准备微调数据(药企NER场景)
training_examples = [
{
"messages": [
{
"role": "system",
"content": "你是药企数据脱敏NER模型。识别文本中的敏感实体,返回JSON列表。"
},
{
"role": "user",
"content": "受试者S015于2024年3月15日退出试验。"
},
{
"role": "assistant",
"content": json.dumps([
{"type": "SUBJECT_ID", "value": "S015", "start": 4, "end": 8},
{"type": "EXACT_DATE", "value": "2024年3月15日", "start": 9, "end": 19}
], ensure_ascii=False)
}
]
},
{
"messages": [
{"role": "system", "content": "你是药企数据脱敏NER模型。识别文本中的敏感实体,返回JSON列表。"},
{"role": "user", "content": "患者张三(SCREEN-042),手机号13900001234,住址上海市浦东新区张江路100号。"},
{"role": "assistant", "content": json.dumps([
{"type": "PATIENT_NAME", "value": "张三", "start": 3, "end": 5},
{"type": "SUBJECT_ID", "value": "SCREEN-042", "start": 6, "end": 17},
{"type": "PHONE_NUMBER", "value": "13900001234", "start": 19, "end": 30},
{"type": "ADDRESS", "value": "上海市浦东新区张江路100号", "start": 32, "end": 47}
], ensure_ascii=False)}
]
}
# ... 准备50-100条类似的标注数据
]
# 保存为JSONL格式(微调数据格式)
with open("ner_training_data.jsonl", "w", encoding="utf-8") as f:
for example in training_examples:
f.write(json.dumps(example, ensure_ascii=False) + "\n")
print("微调数据已准备,共", len(training_examples), "条")
print("使用LLaMA-Factory或swift进行微调:")
print("llamafactory-cli train --model_name_or_path Qwen/Qwen3.5-35B-A3B --dataset ner_training_data.jsonl ...")
⚠️ 踩坑记录
坑1:Qwen3.5通过prompt做NER,通用PII(姓名、身份证、手机号)识别得很好,但是识别不了药企专有的受试者编号格式(S001-S999、SCREEN-XXX)。一开始以为在prompt里写清楚格式就行了,结果测试发现模型经常把S开头的编号跟产品编号混在一起。解决办法是准备50条标注数据做微调,专门标注受试者编号的各种格式。微调用了LLaMA-Factory,2个小时左右就跑完了。微调之后受试者编号的识别率从60%提升到98%。
坑2:Nginx的权限头(X-Data-Scope)跟Dify内部的权限校验必须联动。一开始只配了Nginx层,Dify内部没有做二次校验,结果有人通过改请求头里的X-Department就绕过了权限限制。后来在Dify侧也加了一次权限校验——每次API调用,Dify会读取请求头里的X-Department,然后跟该Agent配置的允许部门列表做比对。两边都校验,Nginx做第一道防线,Dify做兜底。
📎 官方文档:[Nginx](https://nginx.org/en/docs/)
Step 4:审计日志系统(Day 10-11)
在GxP环境里,审计日志不是"可选项",是硬指标。每一次数据访问、每一次数据修改,都要有记录,而且要能导出审计报告给审查员看。
第1、2篇的审计日志是手动记录的(用Python的logging模块写到本地文件),这在单Agent场景下勉强够用。但是到了多Agent跨部门协作的场景,手动记录就完全跟不上了——哪个Agent在什么时间通过什么操作访问了哪个部门的数据?这些信息要自动采集、集中存储、随时可查。
这一步部署EasySearch做审计日志的集中存储和检索。EasySearch是INFINI Labs开源的搜索引擎,兼容Elasticsearch的API,但是部署更轻量。
# 部署EasySearch(Docker方式)
# EasySearch是INFINI Labs开发的国产开源搜索引擎
# 完全兼容Elasticsearch API,Docker镜像很小
docker run -d \
--name easysearch \
--restart always \
-p 9200:9200 \
-p 9300:9300 \
-e "ES_JAVA_OPTS=-Xms4g -Xmx4g" \
-v /data/easysearch/data:/infini/data \
infinilabs/easysearch:latest
# 等待服务启动(大约30秒)
echo "等待EasySearch启动..."
sleep 30
# 验证服务是否正常
curl -s http://localhost:9200/ | python3 -m json.tool
import requests
import json
from datetime import datetime, timedelta
from typing import Optional
EASYSEARCH_URL = "http://localhost:9200"
AUDIT_INDEX = "pharma_audit_log"
# ---- 创建审计日志索引 ----
def create_audit_index():
"""
创建审计日志索引,定义字段映射。
这些字段是GxP审计必须可追溯的。
"""
index_config = {
"settings": {
"number_of_shards": 1,
"number_of_replicas": 0,
# 审计日志保留180天(GxP审查通常要求6个月以上的记录)
"index.lifecycle.name": "audit_retention_180d",
"index.lifecycle.rollover_alias": "audit-write"
},
"mappings": {
"properties": {
"timestamp": {"type": "date", "format": "yyyy-MM-dd'T'HH:mm:ss.SSS+08:00"},
"agent_id": {"type": "keyword"},
"agent_name": {"type": "keyword"},
"department": {"type": "keyword"},
"action": {"type": "keyword"},
"action_detail": {"type": "text", "analyzer": "standard"},
"resource_type": {"type": "keyword"},
"resource_id": {"type": "keyword"},
"target_department": {"type": "keyword"},
"data_scope": {"type": "keyword"},
"user_id": {"type": "keyword"},
"user_name": {"type": "keyword"},
"result": {"type": "keyword"},
"summary": {"type": "text"},
"ip_address": {"type": "ip"},
"session_id": {"type": "keyword"},
# 注意:审计日志只记摘要,不记完整的业务数据
# 避免审计日志本身成为数据泄露源
"data_fingerprint": {"type": "keyword"}
}
}
}
response = requests.put(
f"{EASYSEARCH_URL}/{AUDIT_INDEX}",
headers={"Content-Type": "application/json"},
json=index_config
)
if response.status_code == 200:
print(f"审计日志索引 '{AUDIT_INDEX}' 创建成功")
else:
print(f"创建索引失败:{response.text}")
return response.json()
# ---- 审计日志记录类 ----
class AuditLogger:
"""
审计日志记录器。
所有Agent的数据访问和修改操作都通过这个类记录到EasySearch。
关键设计:只记录操作摘要,不记录完整数据。
比如"读取了临床运营数据"这个操作会记录,但是数据内容不会记录。
"""
def __init__(self, easysearch_url: str = EASYSEARCH_URL):
self.es_url = easysearch_url
self.index = AUDIT_INDEX
def log_access(self, agent_id: str, agent_name: str, department: str,
resource_type: str, resource_id: str,
target_department: str, user_id: str,
ip_address: str, session_id: str,
result: str = "success", summary: str = ""):
"""
记录一次数据访问操作。
"""
log_entry = {
"timestamp": datetime.now().astimezone().isoformat(),
"agent_id": agent_id,
"agent_name": agent_name,
"department": department,
"action": "DATA_ACCESS",
"action_detail": f"访问{resource_type}类型的资源",
"resource_type": resource_type,
"resource_id": resource_id,
"target_department": target_department,
"data_scope": "read",
"user_id": user_id,
"user_name": "", # 由上层填充
"result": result,
"summary": summary,
"ip_address": ip_address,
"session_id": session_id,
"data_fingerprint": "" # 数据指纹,用于追踪数据版本,不存具体内容
}
return self._write_log(log_entry)
def log_modification(self, agent_id: str, agent_name: str, department: str,
resource_type: str, resource_id: str,
modification_type: str, user_id: str,
ip_address: str, session_id: str,
result: str = "success", summary: str = "",
data_fingerprint: str = ""):
"""
记录一次数据修改操作。
modification_type: create/update/delete/status_change
"""
log_entry = {
"timestamp": datetime.now().astimezone().isoformat(),
"agent_id": agent_id,
"agent_name": agent_name,
"department": department,
"action": "DATA_MODIFICATION",
"action_detail": f"{modification_type}操作,资源类型:{resource_type}",
"resource_type": resource_type,
"resource_id": resource_id,
"target_department": department,
"data_scope": "write",
"user_id": user_id,
"user_name": "",
"result": result,
"summary": summary,
"ip_address": ip_address,
"session_id": session_id,
"data_fingerprint": data_fingerprint
}
return self._write_log(log_entry)
def log_cross_department(self, agent_id: str, agent_name: str,
from_department: str, to_department: str,
data_type: str, desensitized: bool,
user_id: str, ip_address: str,
session_id: str, summary: str = ""):
"""
记录一次跨部门数据传递。
这是GxP审查最关注的操作类型。
"""
log_entry = {
"timestamp": datetime.now().astimezone().isoformat(),
"agent_id": agent_id,
"agent_name": agent_name,
"department": from_department,
"action": "CROSS_DEPARTMENT_TRANSFER",
"action_detail": f"从{from_department}向{to_department}传递{data_type}类型数据,脱敏状态:{'已脱敏' if desensitized else '未脱敏'}",
"resource_type": data_type,
"resource_id": "",
"target_department": to_department,
"data_scope": "cross_transfer",
"user_id": user_id,
"user_name": "",
"result": "success" if desensitized else "blocked",
"summary": summary,
"ip_address": ip_address,
"session_id": session_id,
"data_fingerprint": ""
}
# 如果数据未脱敏就尝试跨部门传递,记录为blocked
if not desensitized:
log_entry["result"] = "blocked"
log_entry["action_detail"] += "【安全拦截:未脱敏数据禁止跨部门传递】"
return self._write_log(log_entry)
def _write_log(self, log_entry: dict) -> bool:
"""写入一条审计日志到EasySearch。"""
try:
response = requests.post(
f"{self.es_url}/{self.index}/_doc",
headers={"Content-Type": "application/json"},
json=log_entry
)
return response.status_code in [200, 201]
except Exception as e:
print(f"审计日志写入失败:{e}")
return False
# ---- GxP审计报告导出 ----
def export_gxp_audit_report(start_date: str, end_date: str,
department: Optional[str] = None,
agent_id: Optional[str] = None) -> dict:
"""
导出GxP审计报告。
参数:
- start_date: 起始日期,格式 "2024-01-01"
- end_date: 结束日期,格式 "2024-03-31"
- department: 可选,按部门筛选
- agent_id: 可选,按Agent筛选
返回审计报告的JSON格式数据。
"""
# 构建查询条件
query = {
"query": {
"bool": {
"must": [
{
"range": {
"timestamp": {
"gte": f"{start_date}T00:00:00+08:00",
"lte": f"{end_date}T23:59:59+08:00"
}
}
}
]
}
},
"size": 10000,
"sort": [{"timestamp": {"order": "desc"}}]
}
# 按部门筛选
if department:
query["query"]["bool"]["must"].append({
"term": {"department": department}
})
# 按Agent筛选
if agent_id:
query["query"]["bool"]["must"].append({
"term": {"agent_id": agent_id}
})
# 查询审计日志
response = requests.post(
f"{EASYSEARCH_URL}/{AUDIT_INDEX}/_search",
headers={"Content-Type": "application/json"},
json=query
)
if response.status_code != 200:
print(f"审计日志查询失败:{response.text}")
return {}
results = response.json()
hits = results.get("hits", {}).get("hits", [])
total = results.get("hits", {}).get("total", {}).get("value", 0)
# 生成统计摘要
report = {
"report_metadata": {
"generated_at": datetime.now().isoformat(),
"period_start": start_date,
"period_end": end_date,
"department_filter": department,
"agent_filter": agent_id,
"total_records": total
},
"statistics": {
"total_access_operations": 0,
"total_modification_operations": 0,
"total_cross_department_transfers": 0,
"blocked_operations": 0,
"by_department": {},
"by_agent": {},
"by_action_type": {}
},
"detailed_logs": []
}
for hit in hits:
source = hit["_source"]
report["detailed_logs"].append(source)
# 统计
action = source.get("action", "")
dept = source.get("department", "unknown")
agent = source.get("agent_name", "unknown")
result = source.get("result", "")
if action == "DATA_ACCESS":
report["statistics"]["total_access_operations"] += 1
elif action == "DATA_MODIFICATION":
report["statistics"]["total_modification_operations"] += 1
elif action == "CROSS_DEPARTMENT_TRANSFER":
report["statistics"]["total_cross_department_transfers"] += 1
if result == "blocked":
report["statistics"]["blocked_operations"] += 1
# 按部门统计
report["statistics"]["by_department"][dept] = report["statistics"]["by_department"].get(dept, 0) + 1
# 按Agent统计
report["statistics"]["by_agent"][agent] = report["statistics"]["by_agent"].get(agent, 0) + 1
# 按操作类型统计
report["statistics"]["by_action_type"][action] = report["statistics"]["by_action_type"].get(action, 0) + 1
return report
# ---- 初始化和使用示例 ----
# 创建审计索引
create_audit_index()
# 初始化审计日志记录器
audit_logger = AuditLogger()
# 示例:记录一次数据访问
audit_logger.log_access(
agent_id="clinical_ops_agent",
agent_name="临床运营Agent",
department="clinical_ops",
resource_type="ctms_enrollment_data",
resource_id="trial_C-2024-003",
target_department="clinical_ops",
user_id="user_001",
ip_address="192.168.1.100",
session_id="session_abc123",
summary="查询试验C-2024-003的入组进度"
)
# 示例:记录一次跨部门数据传递(已脱敏)
audit_logger.log_cross_department(
agent_id="clinical_ops_agent",
agent_name="临床运营Agent",
from_department="clinical_ops",
to_department="regulatory",
data_type="enrollment_summary",
desensitized=True,
user_id="user_001",
ip_address="192.168.1.100",
session_id="session_abc123",
summary="向注册部传递试验C-2024-003入组进度汇总(已脱敏)"
)
# 示例:导出GxP审计报告
report = export_gxp_audit_report("2024-01-01", "2024-03-31", department="clinical_ops")
print(f"\nGxP审计报告摘要:")
print(f" 时间范围:{report['report_metadata']['period_start']} 至 {report['report_metadata']['period_end']}")
print(f" 总记录数:{report['report_metadata']['total_records']}")
print(f" 数据访问次数:{report['statistics']['total_access_operations']}")
print(f" 数据修改次数:{report['statistics']['total_modification_operations']}")
print(f" 跨部门传递次数:{report['statistics']['total_cross_department_transfers']}")
print(f" 拦截操作次数:{report['statistics']['blocked_operations']}")
⚠️ 踩坑记录
坑1:EasySearch的默认JVM内存是1GB。在测试环境里审计日志量比较小,没发现问题。但是上线之后,3个Agent同时运行,审计日志量一下子就上来了,EasySearch开始频繁GC,最后直接OOM挂掉。解决办法是在Docker启动参数里把JVM内存调到4GB:-e "ES_JAVA_OPTS=-Xms4g -Xmx4g"。如果你的服务器内存有限,至少也要给2GB。
坑2:一开始审计日志里记录了完整的业务数据(包括临床数据原文),觉得这样"审计更完整"。后来做安全评审的时候发现,审计日志本身就成了数据泄露源——任何能访问审计日志的人都能看到完整的患者信息。改成了只记录操作摘要和数据指纹(用于追踪数据版本),不记录完整数据内容。审计的目的是知道"谁做了什么",不是知道"数据是什么"。
📎 官方文档:[EasySearch - INFINI Labs](https://www.infinilabs.com/)

Step 5:跨Agent工作流编排(Day 11-14)
前面4步把基础设施都搭好了——3个Agent、共享记忆层、脱敏组件、权限网关、审计日志。这一步要把它们串起来,实现跨部门的工作流。
以一个药企里最常见的场景为例:临床方案变更。
临床方案改了一个入组标准,这个变更会影响多个部门:
-
医学部需要评估这个变更对现有文献综述的影响
-
注册部需要检查变更是否符合ICH指南,是否需要提交方案修正
-
PI(主要研究者)和注册负责人需要审批
-
临床运营需要执行变更,更新CTMS数据
-
所有部门需要收到通知
这个工作流不是全自动的。GxP要求关键节点必须有人工审批。所以设计的时候,Agent负责信息汇总和分析,人负责审批和决策。

import requests
import json
from datetime import datetime
from typing import Optional
DIFY_BASE_URL = "http://your-dify-server/v1"
DIFY_API_KEY = "your-dify-workflow-api-key"
# ---- 跨部门工作流定义 ----
# 以"临床方案变更"为例
workflow_definition = {
"name": "临床方案变更协作流程",
"description": "当临床方案发生变更时,自动触发跨部门评估、审批和执行流程",
"trigger": "protocol_amendment",
"nodes": [
{
"id": "step1_impact_assessment",
"name": "医学部影响评估",
"agent": "medical_affairs",
"action": "assess_impact",
"input": {
"protocol_id": "${trigger.protocol_id}",
"amendment_description": "${trigger.amendment_description}",
"changed_fields": "${trigger.changed_fields}"
},
"output": {
"impact_summary": "文本,对现有医学文献和报告的影响分析",
"affected_documents": "列表,受影响的文档清单",
"urgency": "high/medium/low"
},
"timeout_hours": 24,
"description": "医学部Agent评估方案变更对现有医学报告和文献的影响范围"
},
{
"id": "step2_compliance_check",
"name": "注册部合规检查",
"agent": "regulatory",
"action": "compliance_check",
"input": {
"protocol_id": "${trigger.protocol_id}",
"amendment_description": "${trigger.amendment_description}",
"changed_fields": "${trigger.changed_fields}",
"medical_impact": "${step1_impact_assessment.output}"
},
"output": {
"compliance_result": "pass/conditional_pass/fail",
"regulatory_requirements": "需要满足的监管要求清单",
"submission_needed": "是否需要提交方案修正(bool)"
},
"timeout_hours": 24,
"description": "注册部Agent检查方案变更是否符合ICH指南和NMPA要求"
},
{
"id": "step3_human_approval",
"name": "人工审批",
"type": "human_approval",
"approvers": {
"pi": {
"name": "主要研究者",
"notification_method": "email",
"timeout_hours": 48,
"escalation": "department_head"
},
"regulatory_lead": {
"name": "注册负责人",
"notification_method": "email",
"timeout_hours": 48,
"escalation": "regulatory_director"
}
},
"required_approvals": "all",
"input": {
"protocol_id": "${trigger.protocol_id}",
"medical_assessment": "${step1_impact_assessment.output}",
"compliance_check": "${step2_compliance_check.output}",
"amendment_document": "${trigger.amendment_document_url}"
},
"description": "PI和注册负责人共同审批方案变更,需要两人都批准才能继续"
},
{
"id": "step4_execution",
"name": "临床运营执行",
"agent": "clinical_ops",
"action": "execute_amendment",
"input": {
"protocol_id": "${trigger.protocol_id}",
"approval_result": "${step3_human_approval.output}",
"amendment_details": "${trigger.amendment_description}"
},
"output": {
"ctms_updated": "CTMS是否更新完成(bool)",
"sites_notified": "已通知的站点列表",
"enrollment_impact": "对入组的影响评估"
},
"description": "审批通过后,临床运营Agent执行方案变更,更新CTMS和通知各站点"
},
{
"id": "step5_notification",
"name": "全部门通知",
"agent": "clinical_ops",
"action": "broadcast_notification",
"input": {
"protocol_id": "${trigger.protocol_id}",
"amendment_summary": "${trigger.amendment_description}",
"execution_result": "${step4_execution.output}",
"shared_memory_summary": "${shared_summary}"
},
"description": "将方案变更的结果通知所有相关部门,并写入共享记忆层"
}
],
# 节点之间的路由关系
"edges": [
{"from": "step1_impact_assessment", "to": "step2_compliance_check"},
{"from": "step2_compliance_check", "to": "step3_human_approval"},
{"from": "step3_human_approval", "to": "step4_execution", "condition": "approved"},
{"from": "step3_human_approval", "to": "END", "condition": "rejected"},
{"from": "step4_execution", "to": "step5_notification"}
],
# 超时和异常处理
"error_handling": {
"on_timeout": "notify_admin_and_pause",
"on_error": "log_and_retry_once",
"max_retries": 1
}
}
# ---- 工作流执行引擎 ----
class CrossAgentWorkflow:
"""
跨Agent工作流执行引擎。
负责按照工作流定义的顺序,依次调用各个Agent,
处理人工审批节点,管理超时和异常。
"""
def __init__(self, workflow_def: dict):
self.workflow = workflow_def
self.agent_urls = {
"medical_affairs": f"{DIFY_BASE_URL}/chat-messages",
"regulatory": f"{DIFY_BASE_URL}/chat-messages",
"clinical_ops": f"{DIFY_BASE_URL}/chat-messages"
}
self.context = {} # 工作流执行上下文
def execute(self, trigger_data: dict):
"""
启动工作流执行。
trigger_data包含触发工作流的数据,比如方案变更的具体内容。
"""
print(f"\n{'='*60}")
print(f"启动工作流:{self.workflow['name']}")
print(f"触发数据:{json.dumps(trigger_data, ensure_ascii=False)[:200]}...")
print(f"{'='*60}\n")
# 初始化上下文
self.context["trigger"] = trigger_data
self.context["started_at"] = datetime.now().isoformat()
self.context["status"] = "running"
# 按顺序执行节点
for node in self.workflow["nodes"]:
node_id = node["id"]
node_name = node["name"]
print(f"\n--- 执行节点:{node_name} ({node_id}) ---")
# 人工审批节点,单独处理
if node.get("type") == "human_approval":
result = self._handle_human_approval(node)
else:
# Agent节点,调用对应的Agent
result = self._execute_agent_node(node)
# 保存节点结果到上下文
self.context[node_id] = {"output": result, "completed_at": datetime.now().isoformat()}
print(f"节点 {node_name} 执行完成,结果摘要:{str(result)[:150]}...")
# 检查是否需要提前终止(比如审批被拒)
if node.get("type") == "human_approval" and result.get("approved") is False:
print(f"\n工作流提前终止:审批被拒绝")
self.context["status"] = "terminated"
return self.context
self.context["status"] = "completed"
print(f"\n{'='*60}")
print(f"工作流执行完成!")
print(f"{'='*60}")
return self.context
def _execute_agent_node(self, node: dict) -> dict:
"""调用对应Agent执行节点任务。"""
agent_key = node["agent"]
agent_url = self.agent_urls.get(agent_key)
if not agent_url:
return {"error": f"未知的Agent:{agent_key}"}
# 解析输入(从上下文中获取前序节点的结果)
resolved_input = self._resolve_input(node["input"])
# 构建Agent调用请求
prompt = self._build_agent_prompt(node, resolved_input)
payload = {
"inputs": {},
"query": prompt,
"response_mode": "blocking",
"conversation_id": "",
"user": f"workflow_{self.workflow['name']}"
}
try:
response = requests.post(
agent_url,
headers={
"Authorization": f"Bearer {DIFY_API_KEY}",
"Content-Type": "application/json"
},
json=payload,
timeout=300 # Agent节点最多执行5分钟
)
if response.status_code == 200:
result = response.json()
return {"answer": result.get("answer", ""), "status": "success"}
else:
return {"error": response.text, "status": "failed"}
except requests.exceptions.Timeout:
return {"error": "Agent执行超时", "status": "timeout"}
except Exception as e:
return {"error": str(e), "status": "error"}
def _handle_human_approval(self, node: dict) -> dict:
"""
处理人工审批节点。
发送审批通知,等待审批结果。
在实际部署中,这个函数会通过邮件/飞书/企微发送审批通知,
然后提供一个审批链接让审批人在网页上操作。
这里简化为等待用户输入。
"""
approvers = node.get("approvers", {})
print(f"\n⏳ 等待人工审批...")
print(f" 审批人:{', '.join([v['name'] for v in approvers.values()])}")
print(f" 需要审批的内容摘要:{str(node.get('input', {}))[:200]}...")
# 实际部署时的超时和提醒机制
timeout_hours = max(v.get("timeout_hours", 48) for v in approvers.values())
# 模拟审批结果(实际部署时这里是等用户操作)
# 如果超时,自动发送提醒给上级
approval_result = {
"approved": True,
"approver_comments": {
"pi": "同意方案变更,入组标准调整合理",
"regulatory_lead": "合规检查通过,需要提交方案修正至NMPA"
},
"approved_at": datetime.now().isoformat()
}
return approval_result
def _resolve_input(self, input_template: dict) -> dict:
"""
解析节点输入,把 ${xxx} 格式的引用替换为实际值。
"""
resolved = {}
for key, value in input_template.items():
if isinstance(value, str) and value.startswith("${") and value.endswith("}"):
ref_path = value[2:-1]
resolved[key] = self._get_context_value(ref_path)
else:
resolved[key] = value
return resolved
def _get_context_value(self, path: str):
"""从上下文中获取指定路径的值。"""
parts = path.split(".")
current = self.context
for part in parts:
if isinstance(current, dict) and part in current:
current = current[part]
else:
return None
return current
def _build_agent_prompt(self, node: dict, resolved_input: dict) -> str:
"""构建发送给Agent的prompt。"""
prompt = f"【工作流任务 - {node['name']}】\n\n"
prompt += f"任务描述:{node.get('description', '')}\n\n"
prompt += "输入数据:\n"
for key, value in resolved_input.items():
prompt += f" {key}:{json.dumps(value, ensure_ascii=False) if not isinstance(value, str) else value}\n"
prompt += "\n请根据以上信息完成分析,并以JSON格式返回结果。"
return prompt
# ---- 通过Dify API部署工作流 ----
def deploy_workflow_to_dify(workflow_def: dict):
"""
将工作流定义部署到Dify。
Dify的Workflow编排支持条件分支、并行节点、人工审批节点。
"""
# 构建Dify Workflow DSL
dify_workflow = {
"name": workflow_def["name"],
"type": "workflow",
"graph": {
"nodes": [],
"edges": []
}
}
# 转换节点
node_position_x = 100
for node_def in workflow_def["nodes"]:
dify_node = {
"id": node_def["id"],
"type": "agent" if node_def.get("type") != "human_approval" else "approval",
"position": {"x": node_position_x, "y": 200},
"data": {
"title": node_def["name"],
"desc": node_def.get("description", ""),
"agent_app_id": node_def.get("agent", ""),
"timeout": node_def.get("timeout_hours", 24) * 3600
}
}
dify_workflow["graph"]["nodes"].append(dify_node)
node_position_x += 300
# 转换边(路由)
for edge in workflow_def["edges"]:
dify_edge = {
"id": f"{edge['from']}-{edge['to']}",
"source": edge["from"],
"target": edge["to"],
"data": {
"condition": edge.get("condition", "")
}
}
dify_workflow["graph"]["edges"].append(dify_edge)
# 通过Dify API创建Workflow
response = requests.post(
f"{DIFY_BASE_URL}/apps",
headers={
"Authorization": f"Bearer {DIFY_API_KEY}",
"Content-Type": "application/json"
},
json=dify_workflow
)
if response.status_code == 200:
workflow_id = response.json()["id"]
print(f"工作流部署成功!ID:{workflow_id}")
return workflow_id
else:
print(f"工作流部署失败:{response.text}")
return None
# ---- 使用示例:触发"临床方案变更"工作流 ----
# 假设临床方案发生了变更
protocol_amendment_trigger = {
"protocol_id": "C-2024-003",
"amendment_description": "将入组标准中的年龄范围从18-65岁调整为18-75岁,以扩大入组人群覆盖面",
"changed_fields": ["eligibility_criteria.age_range"],
"amendment_version": "v3.1",
"amendment_date": "2024-03-20",
"amendment_document_url": "/documents/protocol/C-2024-003_v3.1.pdf"
}
# 部署工作流
workflow_id = deploy_workflow_to_dify(workflow_definition)
# 执行工作流
if workflow_id:
engine = CrossAgentWorkflow(workflow_definition)
result = engine.execute(protocol_amendment_trigger)
print(f"\n工作流执行结果:")
print(f" 状态:{result['status']}")
print(f" 开始时间:{result['started_at']}")
⚠️ 踩坑记录
坑1:Agent之间传递上下文的时候,信息丢失很严重。比如医学部Agent输出的影响评估有2000字,传到注册部Agent的时候只剩下摘要的200字,因为工作流引擎在传递上下文时做了自动截断。解决办法是在共享记忆里保留"摘要+详情"两层——工作流节点之间传递摘要(短文本,不占上下文窗口),如果某个节点需要详情,通过Memory API去查详情层的完整数据。
坑2:人工审批节点的超时处理。PI出差了48小时没审批,工作流就一直卡在那里等。后面的节点全部阻塞,临床运营的站点都在问"方案到底改不改"。解决办法是给审批节点加超时机制——超过24小时发一次提醒,超过48小时自动升级给PI的上级(部门主任)。如果超过72小时还没审批,工作流进入"暂停"状态,给所有相关人发通知。
📎 官方文档:[Dify Workflow](https://docs.dify.ai/guides/workflow)
4. 运行结果
系统运行了6周之后,做了一次效果统计。对比部署前后的数据:
|
指标 |
部署前 |
部署后 |
变化 |
|---|---|---|---|
|
跨部门数据对接时间 |
3天 |
2-4小时 |
↓83-92% |
|
数据格式转换人工介入 |
80% |
10-15% |
↓81-88% |
|
跨部门数据不一致 |
每月15-20处 |
每月0-2处 |
↓87-100% |
|
方案变更通知延迟 |
2-3天 |
实时(小时级) |
— |
|
审计合规检查通过率 |
82% |
96-98% |
↑14-16% |
|
审计日志完整性 |
手动记录,覆盖率~70% |
自动记录,覆盖率100% |
↑30% |
以上数据基于6周实际运行统计,覆盖3个临床试验项目的跨部门协作场景。测试环境与第1、2篇相同,额外增加1台8C16G服务器运行EasySearch和Nginx网关。
再说说增量资源消耗(在第1、2篇基础上额外需要的):
-
EasySearch:约4GB内存,20GB存储(3个月审计日志量)。这是新增的最大资源消耗项
-
Qwen-NER:因为复用的是已经部署的Qwen3.5,推理本身没有额外成本。微调用了一次性标注成本(50条数据,标注了大约半天),之后没有额外运行时开销
-
Nginx权限网关:几乎不消耗额外资源,就是在现有Nginx上加了几段配置
-
Mem0扩展:摘要层和详情层的向量数据存在已有的PGVector里,增加了约2GB的数据库存储
-
总增量:约1台8C16G服务器,加上微调的一次性人工成本
对比三篇记录的总资源投入:最开始只是一台带GPU的服务器跑Dify+Qwen3.5,到现在加了EasySearch和Nginx网关,整体架构的增长是渐进式的。每加一个组件都有明确的业务需求驱动,不是为了用而用。
5. 总结与下一步
踩坑总结
随着实际的任务推进,这次做下来,积累了不少经验,最后整理5条最核心的:
-
多Agent不是"一个万能Agent"。一开始总想让一个Agent做所有事情,结果prompt写到3000字以上,Agent的指令遵循能力就明显下降了。后来拆成3个Agent,每个Agent的prompt控制在500字以内,效果反而好很多。职责边界必须清晰,宁可多拆不要合并。
-
共享记忆必须脱敏。摘要层+详情层的双层设计,是跨部门数据共享的安全模式。摘要层存脱敏后的事件概要,所有Agent可读;详情层存完整数据,只有本部门的Agent能访问。这个设计虽然增加了一些复杂度,但是在GxP审查的时候特别好用——审查员一看就知道数据是怎么流转的、脱敏是怎么做的。
-
Qwen3.5做NER是意外惊喜。既然已经部署了Qwen3.5做LLM,顺手用它做PII检测,通过prompt的方式做NER,零额外部署成本。后期对精度有更高要求的时候,用50条标注数据微调了一下,识别率从60%提到了98%。这个思路可以推广——已经部署的模型,能复用的能力尽量复用。
-
审计日志不是"可选项"。在GxP环境里,这是硬指标,不是 nice-to-have。EasySearch让审计从"审查前补作业"变成了"日常自动归档"。每次审查的时候,5分钟就能导出完整的审计报告,不用提前两周去各个系统里翻日志。
-
人机协作的边界要画清楚。Agent做信息汇总、影响分析和方案建议,人做审批和决策。关键节点必须有人工介入,这不只是GxP的要求,也是实际需求——方案变更到底要不要执行,这个事情不能让Agent来拍板。
核心洞察
做下来最大的感受是:数据孤岛的本质不是技术问题,是组织问题。 部门之间的数据不通,根源在于各部门的数据格式、权限要求、合规标准都不一样。多Agent系统的价值不是"打破部门墙"——那是组织变革的事——而是"在尊重部门边界的前提下,让必要的数据安全地流动"。
再说一个观点:Agent落地不是从技术出发找问题,是从问题出发搭技术。 这三部曲是按照真实开发顺序记录的,不是"先选好技术栈再找场景",而是"先明确痛点再选方案"。文献整理太慢→用RAG+Agent做自动化。长文档改不动→用Agent+工具链做解析和生成。部门间数据不通→用多Agent+脱敏+审计做协作。每一步都是问题驱动的。
最后,技术选型从来不是最重要的,最重要的是理解业务的颗粒度。 知道CTMS和eTMF的格式差异在哪里,知道受试者编号的格式是什么,知道方案变更的审批流程是什么样的——这些业务知识比选什么向量数据库、用什么Embedding模型重要得多。
本系列基于真实行业痛点,技术方案基于开源工具和通用架构设计,可迁移至真实场景。文中涉及的工具选型仅代表个人实践偏好,不构成商业推荐。如果你有真实的问题需要帮助也可以联系我。
更多推荐



所有评论(0)