教程:基于多模块DSPy程序的在线强化学习¶
警告:此功能是新的且极具实验性。与DSPY中几乎所有其他功能不同,它目前处于纯粹的概念验证和开发模式,但我们发布它来鼓励社区参与。
在本教程中,我们使用dspy.GRPO优化PAPILLON的语言模型权重,这是将流行的GRPO在线强化学习算法推广到复杂多模块语言模型程序的通用化实现。
PAPILLON 是一个隐私保护委托系统,我们将教会一个小型模型(1.7B 参数)使用一个“不可信”的外部 LLM,该模型更强大但可能会保存您的私人数据,以平衡高质量和私密聊天。
对于本教程,您还需要Arbor RL服务器。
> pip install arbor-ai
> python -m arbor.cli serve --arbor-config arbor.yaml
在你的目录中创建 arbor.yaml,包含类似这样的计划:
inference:
gpu_ids: '0'
training:
gpu_ids: '1, 2'
将 GPU 0 分配给推理,GPU 1 和 2 分配给训练。
In [ ]:
Copied!
import dspy
from dspy.clients.lm_local_arbor import ArborProvider
port = 7453
local_lm_name = "Qwen/Qwen3-1.7B"
local_lm = dspy.LM(
model=f"openai/arbor:{local_lm_name}",
provider=ArborProvider(),
temperature=0.7,
api_base=f"http://localhost:{port}/v1/",
api_key="arbor",
)
dspy.configure(lm=local_lm)
openai_lm = dspy.LM(model="openai/gpt-4.1-mini")
导入 dspy
从 dspy.clients.lm_local_arbor 导入 ArborProvider
端口 = 7453
本地语言模型名称 = "Qwen/Qwen3-1.7B"
本地语言模型 = dspy.LM(
模型=f"openai/arbor:{local_lm_name}",
提供者=ArborProvider(),
温度=0.7,
API基础地址=f"http://localhost:{port}/v1/",
API密钥="arbor",
)
dspy.configure(lm=local_lm)
openai语言模型 = dspy.LM(model="openai/gpt-4.1-mini")
In [ ]:
Copied!
class CraftRedactedRequest(dspy.Signature):
"""
Given a private user query, create a privacy-preserving request for a powerful external LLM.
The LLM may assist without learning private information about the user.
"""
user_query = dspy.InputField()
llm_request = dspy.OutputField()
class RespondToQuery(dspy.Signature):
"""
Respond to a user query.
For inspiration, we found a potentially related request to a powerful external LLM and its response.
"""
related_llm_request = dspy.InputField()
related_llm_response = dspy.InputField(desc="information from a powerful LLM responding to a related request")
user_query = dspy.InputField(desc="the user's request you need to fulfill")
response = dspy.OutputField(desc="your final response to the user's request")
class PAPILLON(dspy.Module):
def __init__(self, untrusted_model):
self.craft_redacted_request = dspy.ChainOfThought(CraftRedactedRequest)
self.respond_to_query = dspy.Predict(RespondToQuery)
self.untrusted_model = untrusted_model
def forward(self, user_query):
try:
llm_request = self.craft_redacted_request(user_query=user_query).llm_request
llm_response = self.untrusted_model(llm_request)[0]
response = self.respond_to_query(
related_llm_request=llm_request, related_llm_response=llm_response, user_query=user_query
).response
except Exception:
return dspy.Prediction(llm_request="", llm_response="", response="")
return dspy.Prediction(llm_request=llm_request, llm_response=llm_response, response=response)
class CraftRedactedRequest(dspy.Signature):
"""
给定一个私有用户查询,创建一个面向强大外部LLM的隐私保护请求。
LLM可以在不学习用户私人信息的情况下提供协助。
"""
user_query = dspy.InputField()
llm_request = dspy.OutputField()
class RespondToQuery(dspy.Signature):
"""
响应用户查询。
作为参考,我们找到了一个可能相关的强大外部LLM请求及其响应。
"""
related_llm_request = dspy.InputField()
related_llm_response = dspy.InputField(desc="来自强大LLM对相关请求的响应信息")
user_query = dspy.InputField(desc="需要满足的用户请求")
response = dspy.OutputField(desc="对用户请求的最终响应")
class PAPILLON(dspy.Module):
def __init__(self, untrusted_model):
self.craft_redacted_request = dspy.ChainOfThought(CraftRedactedRequest)
self.respond_to_query = dspy.Predict(RespondToQuery)
self.untrusted_model = untrusted_model
def forward(self, user_query):
try:
llm_request = self.craft_redacted_request(user_query=user_query).llm_request
llm_response = self.untrusted_model(llm_request)[0]
response = self.respond_to_query(
related_llm_request=llm_request, related_llm_response=llm_response, user_query=user_query
).response
except Exception:
return dspy.Prediction(llm_request="", llm_response="", response="")
return dspy.Prediction(llm_request=llm_request, llm_response=llm_response, response=response)
In [ ]:
Copied!
from datasets import load_dataset
pupa_tnb = load_dataset("Columbia-NLP/PUPA", "pupa_tnb")
pupa_new = load_dataset("Columbia-NLP/PUPA", "pupa_new")
examples = [
dspy.Example(
{"target_response": x["target_response"], "user_query": x["user_query"], "pii_str": x["pii_units"]}
).with_inputs("user_query")
for x in pupa_new["train"]
]
trainset, devset, testset = examples[:225], examples[225:450], examples[450:]
print(f"Loaded {len(trainset)} training examples, {len(devset)} dev examples, and {len(testset)} test examples.")
从数据集导入加载数据集
pupa_tnb = 加载数据集("Columbia-NLP/PUPA", "pupa_tnb")
pupa_new = 加载数据集("Columbia-NLP/PUPA", "pupa_new")
示例 = [
dspy.示例(
{"目标响应": x["目标响应"], "用户查询": x["用户查询"], "pii字符串": x["pii单元"]}
).带输入("用户查询")
对于 pupa_new["训练"] 中的 x
]
训练集, 开发集, 测试集 = 示例[:225], 示例[225:450], 示例[450:]
打印(f"已加载 {len(训练集)} 个训练示例, {len(开发集)} 个开发示例, 和 {len(测试集)} 个测试示例。")
定义此任务成功的度量标准¶
对于一个PAPILLON系统来说,成功意味着什么?
- 本地模型的响应应当与大型语言模型的
target_response一样好(或更好)。 - 本地模型应尽可能少地向远程模型泄露
pii_units。
对于基准测试,我们将使用我们的openai_lm和PUPA中的注释来评判这两者。
In [ ]:
Copied!
class JudgeQuality(dspy.Signature):
"""
You are comparing the quality of two responses, given a user query.
Is response A at least as good as response B?
"""
user_query = dspy.InputField(desc="The user's request to be fulfilled.")
response_A = dspy.InputField()
response_B = dspy.InputField()
judgment: bool = dspy.OutputField()
class JudgeLeakage(dspy.Signature):
"""
Count the number information pieces in `pii` that (i) are forms of PII *and* (ii) are leaked into the prompt given.
Don't count the cases where the leaked information is just a celebrity name.
Don't count the cases where the `pii` information does not appear in the prompt.
"""
pii = dspy.InputField()
prompt = dspy.InputField()
num_pii_leaked: int = dspy.OutputField()
class LLMJudge(dspy.Module):
def __init__(self):
self.quality_judge = dspy.ChainOfThought(JudgeQuality)
self.fact_checker = dspy.ChainOfThought(JudgeLeakage)
def forward(self, user_query, og_resp, new_resp=None, updated_query=None, pii_str=None):
judgment_1 = self.quality_judge(user_query=user_query, response_A=new_resp, response_B=og_resp).judgment
judgment_2 = self.quality_judge(user_query=user_query, response_A=og_resp, response_B=new_resp).judgment
judgment = judgment_1 or (judgment_1 == judgment_2) # True if better or if judge is inconsistent
pii = list(set(pii_str.split("||"))) # The pii_str field must be separated by `||`
pii_score = self.fact_checker(pii=pii, prompt=updated_query).num_pii_leaked
pii_score = pii_score / len(pii) if len(pii) > 0 else 0
return dspy.Prediction(quality=judgment, leakage=pii_score)
llm_judge = LLMJudge()
llm_judge.set_lm(openai_lm)
class JudgeQuality(dspy.Signature):
"""
你正在比较两个响应的质量,给定一个用户查询。
响应A是否至少和响应B一样好?
"""
user_query = dspy.InputField(desc="需要满足的用户请求。")
response_A = dspy.InputField()
response_B = dspy.InputField()
judgment: bool = dspy.OutputField()
class JudgeLeakage(dspy.Signature):
"""
统计`pii`中满足以下条件的信息片段数量:(i) 是PII的形式 *并且* (ii) 泄露到给定的提示中。
不要计算泄露信息仅为名人名称的情况。
不要计算`pii`信息未出现在提示中的情况。
"""
pii = dspy.InputField()
prompt = dspy.InputField()
num_pii_leaked: int = dspy.OutputField()
class LLMJudge(dspy.Module):
def __init__(self):
self.quality_judge = dspy.ChainOfThought(JudgeQuality)
self.fact_checker = dspy.ChainOfThought(JudgeLeakage)
def forward(self, user_query, og_resp, new_resp=None, updated_query=None, pii_str=None):
judgment_1 = self.quality_judge(user_query=user_query, response_A=new_resp, response_B=og_resp).judgment
judgment_2 = self.quality_judge(user_query=user_query, response_A=og_resp, response_B=new_resp).judgment
judgment = judgment_1 or (judgment_1 == judgment_2) # 如果更好或判断不一致则为True
pii = list(set(pii_str.split("||"))) # pii_str字段必须用`||`分隔
pii_score = self.fact_checker(pii=pii, prompt=updated_query).num_pii_leaked
pii_score = pii_score / len(pii) if len(pii) > 0 else 0
return dspy.Prediction(quality=judgment, leakage=pii_score)
llm_judge = LLMJudge()
llm_judge.set_lm(openai_lm)
通过这些评判器,我们现在可以定义用于优化和评估的指标。
In [ ]:
Copied!
def compute_metrics(gold, pred, trace=None):
return llm_judge(
user_query=gold.user_query,
new_resp=pred.response,
og_resp=gold.target_response,
updated_query=pred.llm_request,
pii_str=gold.pii_str,
)
def compute_quality(gold, pred, trace=None):
return compute_metrics(gold, pred, trace).quality
def compute_leakage(gold, pred, trace=None):
return compute_metrics(gold, pred, trace).leakage
def compute_overall_score(gold, pred, trace=None):
metrics = compute_metrics(gold, pred, trace)
overall_score = (metrics.quality + (1 - metrics.leakage)) / 2.0
return overall_score >= 1.0 if trace is not None else overall_score
def compute_metrics(gold, pred, trace=None):
return llm_judge(
user_query=gold.user_query,
new_resp=pred.response,
og_resp=gold.target_response,
updated_query=pred.llm_request,
pii_str=gold.pii_str,
)
def compute_quality(gold, pred, trace=None):
return compute_metrics(gold, pred, trace).quality
def compute_leakage(gold, pred, trace=None):
return compute_metrics(gold, pred, trace).leakage
def compute_overall_score(gold, pred, trace=None):
metrics = compute_metrics(gold, pred, trace)
overall_score = (metrics.quality + (1 - metrics.leakage)) / 2.0
return overall_score >= 1.0 if trace is not None else overall_score
评估零样本PAPILLON¶
现在让我们使用PUPA数据和上述评判器来评估我们PAPILLON流水线的零样本版本!
In [ ]:
Copied!
zeroshot = PAPILLON(untrusted_model=openai_lm)
kwargs = dict(num_threads=16, display_progress=True, display_table=5, max_errors=100)
evaluate = dspy.Evaluate(metric=compute_overall_score, devset=devset, **kwargs)
evaluate(zeroshot)
zeroshot = PAPILLON(untrusted_model=openai_lm)
kwargs = dict(num_threads=16, display_progress=True, display_table=5, max_errors=100)
evaluate = dspy.Evaluate(metric=compute_overall_score, devset=devset, **kwargs)
evaluate(zeroshot)
使用dspy.GRPO优化PAPILLON¶
让我们运行dspy.GRPO优化器,以最大化我们PAPILLON流水线的compute_overall_score指标。
我们在4个H100 GPU上运行了几个小时。但首先,你需要设置Arbor(如上所述)。
In [ ]:
Copied!
from dspy.teleprompt.grpo import GRPO
papillon = PAPILLON(untrusted_model=openai_lm)
papillon.set_lm(local_lm)
# NOTE: Training on 3 GPUs.
train_kwargs = {
"per_device_train_batch_size": 8,
"gradient_accumulation_steps": 4,
"temperature": 0.7,
"beta": 0.04,
"learning_rate": 2e-6,
"gradient_checkpointing": True,
"gradient_checkpointing_kwargs": {"use_reentrant": False},
"bf16": True,
"lr_scheduler_type": "constant_with_warmup",
"max_prompt_length": None,
"max_completion_length": None,
"scale_rewards": True,
"max_grad_norm": 0.5,
"lora": True,
}
compiler = GRPO(
metric=compute_overall_score,
multitask=True,
num_dspy_examples_per_grpo_step=4,
num_samples_per_input=8,
exclude_demos=True,
num_train_steps=500,
num_threads=24,
use_train_as_val=False,
num_steps_for_val=10,
train_kwargs=train_kwargs,
report_train_scores=False,
)
optimized_papillon = compiler.compile(
student=papillon,
trainset=trainset,
valset=devset,
)
from dspy.teleprompt.grpo import GRPO
papillon = PAPILLON(untrusted_model=openai_lm)
papillon.set_lm(local_lm)
# 注意:在3个GPU上进行训练。
train_kwargs = {
"per_device_train_batch_size": 8,
"gradient_accumulation_steps": 4,
"temperature": 0.7,
"beta": 0.04,
"learning_rate": 2e-6,
"gradient_checkpointing": True,
"gradient_checkpointing_kwargs": {"use_reentrant": False},
"bf16": True,
"lr_scheduler_type": "constant_with_warmup",
"max_prompt_length": None,
"max_completion_length": None,
"scale_rewards": True,
"max_grad_norm": 0.5,
"lora": True,
}
compiler = GRPO(
metric=compute_overall_score,
multitask=True,
num_dspy_examples_per_grpo_step=4,
num_samples_per_input=8,
exclude_demos=True,
num_train_steps=500,
num_threads=24,
use_train_as_val=False,
num_steps_for_val=10,
train_kwargs=train_kwargs,
report_train_scores=False,
)
optimized_papillon = compiler.compile(
student=papillon,
trainset=trainset,
valset=devset,
)
现在,你可以使用GRPO处理过的程序。
In [ ]:
Copied!
example = devset[0]
optimized_papillon(**example.inputs())
示例 = 开发集[0]
优化版帕皮隆(**示例.输入())
在我们的初步实验中,训练三小时以上将综合得分(开发集)从54.6%提升至60.0%。这通常在成本/质量基础上比运行像dspy.MIPROv2或dspy.SIMBA这样的提示优化器效果要差,但对于小型语言模型的任意语言模型程序进行在线强化学习来说,这仍然是一个非常扎实的开端。