八:手撕代码与项目深挖
实现一个将对话数据转换为训练格式的Data Collator,包含loss masking¶
这个 Data Collator 接收一个 batch 的对话样本,每个样本包含 input_ids 和 labels(已经对 prompt 部分设置为 -100 进行掩码)。它的主要职责是动态填充同一 batch 内的序列到相同长度,并对 labels 做相应的填充(填充值为 -100),以保证模型在计算损失时忽略填充部分和 prompt 部分。
import torch
from typing import List, Dict
class SFTDataCollator:
"""
数据整理器,用于监督微调(SFT)。
期望每个样本是一个字典,包含 'input_ids' 和 'labels',
其中 labels 中已经对 prompt 部分做了 loss masking(设为 -100)。
"""
def __init__(self, pad_token_id: int = 0):
self.pad_token_id = pad_token_id
def __call__(self, features: List[Dict[str, torch.Tensor]]) -> Dict[str, torch.Tensor]:
# 1. 找到 batch 中最长的序列长度
max_len = max(f['input_ids'].size(0) for f in features)
padded_input_ids = []
padded_labels = []
attention_masks = []
for f in features:
input_ids = f['input_ids']
labels = f['labels']
seq_len = input_ids.size(0)
pad_len = max_len - seq_len
# 2. 对 input_ids 右侧填充 pad_token_id
padded_input = torch.cat([
input_ids,
torch.full((pad_len,), self.pad_token_id, dtype=input_ids.dtype)
])
# 3. 对 labels 右侧填充 -100(损失忽略)
padded_label = torch.cat([
labels,
torch.full((pad_len,), -100, dtype=labels.dtype)
])
# 4. 注意力掩码:真实 token 为 1,填充为 0
attn_mask = torch.cat([
torch.ones(seq_len, dtype=torch.long),
torch.zeros(pad_len, dtype=torch.long)
])
padded_input_ids.append(padded_input)
padded_labels.append(padded_label)
attention_masks.append(attn_mask)
# 5. 堆叠成 batch 张量
batch = {
'input_ids': torch.stack(padded_input_ids),
'labels': torch.stack(padded_labels),
'attention_mask': torch.stack(attention_masks)
}
return batch
说明:
-
为了支持 loss masking,我们在预处理阶段就已经将 prompt 对应的 labels 设置为 -100,Data Collator 只需在填充时对填充位置同样置 -100,保证这些位置不参与损失计算。
-
attention_mask用于让模型忽略填充位置的信息。 -
这种设计简洁可靠,符合 HuggingFace Trainer 的接口要求。
手写代码将Alpaca格式数据转为ChatML模板¶
Alpaca 格式通常包含 instruction、input(可选)和 output 字段。ChatML 模板则采用 <|im_start|> 和 <|im_end|> 标记来区分角色。转换函数需将指令和回答包装成 ChatML 对话格式。
def alpaca_to_chatml(instruction: str, input_text: str = "", output: str = "") -> str:
"""
将 Alpaca 格式的单轮数据转换为 ChatML 格式字符串。
如果 input_text 为空,则只使用 instruction;否则将 instruction 和 input 拼接。
"""
# 构建用户消息内容
if input_text:
user_content = f"{instruction}\n{input_text}"
else:
user_content = instruction
# 组装 ChatML
chatml = (
f"<|im_start|>user\n{user_content}<|im_end|>\n"
f"<|im_start|>assistant\n{output}<|im_end|>"
)
return chatml
示例:
instruction = "写一首关于秋天的诗。"
output = "秋风起,叶落知多少。"
chatml_str = alpaca_to_chatml(instruction, output=output)
print(chatml_str)
# <|im_start|>user
# 写一首关于秋天的诗。<|im_end|>
# <|im_start|>assistant
# 秋风起,叶落知多少。<|im_end|>
说明:
-
为保持一致性,通常会在 tokenizer 中添加
<|im_start|>和<|im_end|>为特殊 token。 -
该函数仅处理单轮对话;多轮对话可按类似逻辑循环拼接。
使用PyTorch实现Multi-head Self-Attention(带mask)¶
实现一个标准的多头自注意力模块,支持因果掩码(causal mask)和填充掩码(padding mask)。为简洁,这里假设输入是 (batch, seq_len, d_model),且已通过线性层得到 Q, K, V。
import torch
import torch.nn as nn
import torch.nn.functional as F
class MultiHeadSelfAttention(nn.Module):
def __init__(self, d_model: int, n_heads: int, dropout: float = 0.1):
super().__init__()
assert d_model % n_heads == 0
self.d_model = d_model
self.n_heads = n_heads
self.d_k = d_model // n_heads
self.W_q = nn.Linear(d_model, d_model)
self.W_k = nn.Linear(d_model, d_model)
self.W_v = nn.Linear(d_model, d_model)
self.W_o = nn.Linear(d_model, d_model)
self.dropout = nn.Dropout(dropout)
def forward(self, x: torch.Tensor, mask: torch.Tensor = None) -> torch.Tensor:
B, T, C = x.shape
# 线性变换并拆分为多头
q = self.W_q(x).view(B, T, self.n_heads, self.d_k).transpose(1, 2) # (B, n_heads, T, d_k)
k = self.W_k(x).view(B, T, self.n_heads, self.d_k).transpose(1, 2)
v = self.W_v(x).view(B, T, self.n_heads, self.d_k).transpose(1, 2)
# 缩放点积注意力
attn_scores = torch.matmul(q, k.transpose(-2, -1)) / (self.d_k ** 0.5) # (B, n_heads, T, T)
# 应用掩码
if mask is not None:
# mask 形状通常为 (B, 1, 1, T) 或 (B, 1, T, T),这里兼容两种
attn_scores = attn_scores.masked_fill(mask == 0, float('-inf'))
attn_weights = F.softmax(attn_scores, dim=-1)
attn_weights = self.dropout(attn_weights)
# 加权求和
out = torch.matmul(attn_weights, v) # (B, n_heads, T, d_k)
out = out.transpose(1, 2).contiguous().view(B, T, C) # (B, T, d_model)
out = self.W_o(out)
return out
说明:
-
掩码
mask通常来自attention_mask,其中 0 表示需要屏蔽的位置。在实际使用中,可以通过一个函数生成因果掩码并与填充掩码结合。 -
该实现省略了单独的 bias 和缓存机制,专注于注意力核心。
编写代码实现RoPE位置编码的apply_rotary_emb函数¶
RoPE(旋转位置编码)通过对查询和键向量施加旋转变换注入位置信息。下面实现一个通用的应用函数。
import torch
def apply_rotary_emb(x: torch.Tensor, cos: torch.Tensor, sin: torch.Tensor) -> torch.Tensor:
"""
对输入张量 x 应用旋转位置编码。
参数:
x: 形状为 (..., seq_len, dim),其中 dim 是偶数。
cos: 形状为 (seq_len, dim//2) 或 (..., seq_len, dim//2)
sin: 形状同上
返回:
旋转后的张量,形状与 x 相同。
"""
# 将 x 按最后一维拆分为两半,分别旋转
dim = x.shape[-1]
x_rot = x[..., :dim//2] # 前半部分
x_pass = x[..., dim//2:] # 后半部分(如果存在不对称处理,通常 RoPE 是整维度)
# 实际 RoPE 将整个向量分成对,每对 (x0, x1) 旋转角度 θ,
# 但更高效的方式是将向量视为复数,利用复旋转。
# 这里实现常用的复数旋转形式。
x_reshaped = x.view(*x.shape[:-1], -1, 2)
x_cos = x_reshaped * cos.unsqueeze(-2) # 伸缩 cos 部分
# 构建旋转后的另一半:交换两维并取负
x_sin = torch.stack([-x_reshaped[..., 1], x_reshaped[..., 0]], dim=-1)
x_sin = x_sin * sin.unsqueeze(-2)
rotated = (x_cos + x_sin).view_as(x)
return rotated
def precompute_rotary_embeddings(seq_len: int, dim: int, theta: float = 10000.0):
"""预计算旋转位置编码的 cos 和 sin 表"""
# 频率
freqs = 1.0 / (theta ** (torch.arange(0, dim, 2).float() / dim)) # (dim/2)
t = torch.arange(seq_len).float()
freqs = torch.outer(t, freqs) # (seq_len, dim/2)
cos = torch.cos(freqs)
sin = torch.sin(freqs)
return cos, sin
说明:
-
上述实现采用了复数旋转的观点:将向量视为成对的 (a, b),旋转角度 θ 后变为 (a cosθ - b sinθ, a sinθ + b cosθ)。
-
实际使用时,
apply_rotary_emb分别对 query 和 key 张量应用 cos/sin,且通常只对头部维度的一半或特定维度施加 RoPE。
实现一个简化版GPT模型的forward,包含transformer block和lm_head¶
下面实现一个简易的 Transformer 解码器层和完整的 GPT 模型前向传播,包括词嵌入、位置编码(这里简单用可学习位置编码)、多层 transformer block 和 lm_head。
import torch.nn as nn
import torch
class TransformerBlock(nn.Module):
def __init__(self, d_model: int, n_heads: int, d_ff: int, dropout: float = 0.1):
super().__init__()
self.attn = MultiHeadSelfAttention(d_model, n_heads, dropout)
self.ln1 = nn.LayerNorm(d_model)
self.ffn = nn.Sequential(
nn.Linear(d_model, d_ff),
nn.GELU(),
nn.Linear(d_ff, d_model),
nn.Dropout(dropout)
)
self.ln2 = nn.LayerNorm(d_model)
self.dropout = nn.Dropout(dropout)
def forward(self, x, mask=None):
x = x + self.dropout(self.attn(self.ln1(x), mask))
x = x + self.dropout(self.ffn(self.ln2(x)))
return x
class SimpleGPT(nn.Module):
def __init__(self, vocab_size, d_model=512, n_heads=8, n_layers=6, d_ff=2048, max_len=1024, dropout=0.1):
super().__init__()
self.token_embed = nn.Embedding(vocab_size, d_model)
self.pos_embed = nn.Embedding(max_len, d_model)
self.blocks = nn.ModuleList([TransformerBlock(d_model, n_heads, d_ff, dropout) for _ in range(n_layers)])
self.ln_f = nn.LayerNorm(d_model)
self.lm_head = nn.Linear(d_model, vocab_size, bias=False)
self.max_len = max_len
self.dropout = nn.Dropout(dropout)
def forward(self, input_ids, attention_mask=None):
B, T = input_ids.shape
# 构造因果掩码
causal_mask = torch.tril(torch.ones(T, T, device=input_ids.device)).view(1, 1, T, T)
if attention_mask is not None:
# attention_mask 形状 (B, T),需要扩展并合并
mask = attention_mask[:, None, None, :] * causal_mask
else:
mask = causal_mask
# Token Embedding + Position Embedding
positions = torch.arange(0, T, device=input_ids.device).unsqueeze(0).expand(B, -1)
tok_emb = self.token_embed(input_ids)
pos_emb = self.pos_embed(positions)
x = self.dropout(tok_emb + pos_emb)
# 通过 transformer blocks
for block in self.blocks:
x = block(x, mask)
x = self.ln_f(x)
logits = self.lm_head(x) # (B, T, vocab_size)
return logits
说明:
-
该简化版模型省略了权重复用等细节,但清晰展示了整体流程。
-
因果掩码确保每个位置只能注意到它前面的位置。
-
lm_head输出 logits,可直接用于交叉熵损失计算。
实现LoRA线性层的定义及参数初始化¶
LoRA 在原始线性层旁添加低秩矩阵 A 和 B,前向时输出为 Wx + (α/r) * BAx。下面实现一个可配置的 LoRA 线性层。
import torch
import torch.nn as nn
import math
class LoRALinear(nn.Module):
def __init__(self, in_features: int, out_features: int, r: int = 8, lora_alpha: int = 16,
bias: bool = True, dropout: float = 0.0):
super().__init__()
self.in_features = in_features
self.out_features = out_features
self.r = r
self.lora_alpha = lora_alpha
self.scaling = lora_alpha / r
# 原始线性层,权重冻结
self.linear = nn.Linear(in_features, out_features, bias=bias)
self.linear.weight.requires_grad = False
if bias:
self.linear.bias.requires_grad = False
# LoRA 低秩矩阵 A (降维) 和 B (升维)
self.lora_A = nn.Parameter(torch.zeros(r, in_features))
self.lora_B = nn.Parameter(torch.zeros(out_features, r))
self.lora_dropout = nn.Dropout(dropout) if dropout > 0 else nn.Identity()
# 参数初始化
self.reset_lora_parameters()
def reset_lora_parameters(self):
# A 使用 Kaiming 初始化或正态分布,B 初始化为零,确保初始时 BAx = 0
nn.init.kaiming_uniform_(self.lora_A, a=math.sqrt(5))
nn.init.zeros_(self.lora_B)
def forward(self, x):
# 原始线性输出 (冻结)
result = self.linear(x)
# LoRA 增量
lora_out = self.lora_dropout(x) @ self.lora_A.T # (..., in) @ (in, r) -> (..., r)
lora_out = lora_out @ self.lora_B.T # (..., r) @ (r, out) -> (..., out)
result = result + self.scaling * lora_out
return result
说明:
-
原始线性层权重被冻结 (
requires_grad=False)。 -
LoRA 的 A 矩阵维度为
(r, in),B 为(out, r)。scaling因子为lora_alpha / r。 -
通过将 B 初始化为零,训练开始时 LoRA 的增量为零,不会破坏原始输出。
编写函数将LoRA权重合并到原模型权重中¶
训练后,可以将 LoRA 矩阵与原始权重合并,生成一个等效的稠密线性层,便于推理部署。
def merge_lora_to_linear(lora_layer: LoRALinear) -> nn.Linear:
"""将 LoRA 权重合并到原始线性层中,返回新的合并后的线性层"""
# 确保原权重与 LoRA 在同一设备
W = lora_layer.linear.weight.data # (out, in)
A = lora_layer.lora_A.data # (r, in)
B = lora_layer.lora_B.data # (out, r)
# 计算增量 ΔW = (lora_alpha / r) * B @ A
delta_W = (lora_layer.lora_alpha / lora_layer.r) * torch.matmul(B, A) # (out, in)
merged_weight = W + delta_W
# 构建新的线性层
new_linear = nn.Linear(lora_layer.in_features, lora_layer.out_features,
bias=lora_layer.linear.bias is not None)
new_linear.weight.data = merged_weight
if lora_layer.linear.bias is not None:
new_linear.bias.data = lora_layer.linear.bias.data
return new_linear
说明:
-
合并后即可删除 LoRA 参数,仅保留等价的标准线性层,实现零额外推理开销。
-
一般对整个模型所有 LoRA 层执行此操作,然后保存为新的模型文件。
手写Top-P(nucleus)采样算法¶
Top-P 采样从累积概率超过阈值 p 的最小 token 集合中随机采样下一个 token。
import torch
def top_p_sampling(logits: torch.Tensor, p: float = 0.9, temperature: float = 1.0) -> int:
"""
对 logits 应用 Top-P (nucleus) 采样,返回采样得到的 token id。
参数:
logits: 形状 (vocab_size,) 的未归一化 logits
p: 核采样阈值
temperature: 温度系数,控制分布锐度
"""
# 1. 温度缩放
if temperature > 0:
logits = logits / temperature
# 2. 转为概率
probs = torch.softmax(logits, dim=-1)
# 3. 按概率降序排列
sorted_probs, sorted_indices = torch.sort(probs, descending=True)
cumsum_probs = torch.cumsum(sorted_probs, dim=-1)
# 4. 去除累积概率超过 p 的尾部
mask = cumsum_probs <= p
# 至少保留一个 token,所以将第一个 token 强制保留
mask[0] = True
# 5. 对筛选后的子集重新归一化
filtered_probs = sorted_probs * mask.float()
filtered_probs = filtered_probs / filtered_probs.sum()
# 6. 从筛选后的分布中采样
sampled_idx = torch.multinomial(filtered_probs, 1).item()
# 7. 映射回原始 token id
return sorted_indices[sampled_idx].item()
说明:
-
p通常设为 0.9 或 0.95,以避免过于离奇的 token。 -
温度
temperature可调节随机性:越高越随机,越低越贪婪。
手写Top-K采样算法¶
Top-K 采样仅从概率最高的 K 个 token 中采样。
def top_k_sampling(logits: torch.Tensor, k: int = 50, temperature: float = 1.0) -> int:
"""
对 logits 应用 Top-K 采样,返回采样得到的 token id。
参数:
logits: 形状 (vocab_size,) 的未归一化 logits
k: 保留的最高概率 token 数量
temperature: 温度系数
"""
if temperature > 0:
logits = logits / temperature
# 1. 计算概率
probs = torch.softmax(logits, dim=-1)
# 2. 选取概率最高的 k 个 token
topk_probs, topk_indices = torch.topk(probs, k, dim=-1)
# 3. 重新归一化
topk_probs = topk_probs / topk_probs.sum()
# 4. 采样
sampled_idx = torch.multinomial(topk_probs, 1).item()
return topk_indices[sampled_idx].item()
说明:
-
k通常取 50、100 等,太小会丧失多样性,太大则可能包含低概率的噪声。 -
实际应用中常将 Top-K 和 Top-P 结合,先 Top-K 过滤,再应用 Top-P。
以上代码覆盖了从数据处理到模型结构、位置编码、参数高效微调以及解码策略的关键实现,可作为大模型微调和推理的基础组件。
实现KV Cache的更新逻辑(单层,考虑batch)¶
在自回归推理中,KV Cache 存储了已生成 token 的 Key 和 Value 张量。每次新 token 到来时,只需计算该 token 的 K, V,并与之前的缓存拼接。以下实现了一个单层 Transformer 的 KV Cache 更新逻辑,适用于批量推理。
import torch
def update_kv_cache(
new_k: torch.Tensor, # [batch, num_heads, 1, head_dim] 新 token 的 Key
new_v: torch.Tensor, # [batch, num_heads, 1, head_dim] 新 token 的 Value
cache_k: torch.Tensor, # [batch, num_heads, seq_len, head_dim] 或 None
cache_v: torch.Tensor,
max_cache_len: int = None # 最大缓存长度,超出则截断
):
"""
更新 KV Cache:将新 token 的 K,V 拼接到历史缓存上。
如果 cache 为空,则直接返回新的 K,V(增加 seq_len 维度)。
返回更新后的 cache_k, cache_v。
"""
if cache_k is None:
# 初始化缓存
return new_k, new_v
# 沿 seq_len 维度拼接 (dim=2)
updated_k = torch.cat([cache_k, new_k], dim=2)
updated_v = torch.cat([cache_v, new_v], dim=2)
# 若设置了最大长度,则仅保留最近的 max_cache_len 个 token
if max_cache_len is not None and updated_k.size(2) > max_cache_len:
updated_k = updated_k[:, :, -max_cache_len:, :]
updated_v = updated_v[:, :, -max_cache_len:, :]
return updated_k, updated_v
说明:
-
新 token 的 K,V 形状应为
(batch, num_heads, 1, head_dim),通过unsqueeze添加长度维度。 -
第一次调用时
cache_k=None,函数直接返回新 K,V 作为初始缓存。 -
max_cache_len可用于限制缓存长度,避免显存溢出。
编写自定义Trainer回调函数,在SFT每个epoch结束时进行生成测试¶
我们可以通过继承 TrainerCallback 来实现一个回调,在每个 epoch 结束时用几个固定的 prompt 让模型生成回答,并记录日志。
from transformers import TrainerCallback, TrainerControl, TrainerState, TrainingArguments
import torch
class GenerationEvalCallback(TrainerCallback):
"""每个 epoch 结束时,使用测试 prompt 生成回答,打印或记录日志"""
def __init__(self, test_prompts: list, tokenizer, max_new_tokens=50):
self.test_prompts = test_prompts
self.tokenizer = tokenizer
self.max_new_tokens = max_new_tokens
def on_epoch_end(self, args: TrainingArguments, state: TrainerState, control: TrainerControl, **kwargs):
model = kwargs['model']
model.eval()
print(f"\n=== Epoch {state.epoch:.2f} 生成测试 ===")
for prompt in self.test_prompts:
inputs = self.tokenizer(prompt, return_tensors='pt').to(model.device)
with torch.no_grad():
outputs = model.generate(
**inputs,
max_new_tokens=self.max_new_tokens,
do_sample=False,
pad_token_id=self.tokenizer.eos_token_id
)
generated = self.tokenizer.decode(outputs[0], skip_special_tokens=True)
print(f"Prompt: {prompt}\nGenerated: {generated}\n")
model.train() # 恢复训练模式
使用方式:
trainer.add_callback(GenerationEvalCallback(
test_prompts=["你好,请介绍一下人工智能。", "写一首关于春天的诗。"],
tokenizer=tokenizer
))
实现一个序列打包(packing)函数,将多个输入拼接到一个序列并生成attention mask¶
该函数将多个样本的 input_ids 拼接成一个长序列,同时构建分块对角注意力掩码,使得不同样本之间不能互相关注。
import torch
def pack_sequences(
input_ids_list: list[torch.Tensor], # 每个样本的 input_ids (1D)
labels_list: list[torch.Tensor], # 每个样本的 labels
pad_token_id: int = 0,
eos_token_id: int = 2
) -> dict:
"""
将多个样本拼接为一个 packed 序列,返回 input_ids, labels, attention_mask。
attention_mask 采用 4D 分块对角形式:形状 (1, 1, total_len, total_len)
"""
# 在样本之间插入 EOS 作为分隔(可选)
packed_input_ids = []
packed_labels = []
segment_boundaries = [] # 每个样本的 (start, end) 索引(包含)
offset = 0
for inp, lab in zip(input_ids_list, labels_list):
L = inp.size(0)
packed_input_ids.append(inp)
packed_labels.append(lab)
segment_boundaries.append((offset, offset + L - 1))
offset += L
# 插入分隔 token (EOS),其 label 设为 -100(忽略)
packed_input_ids.append(torch.tensor([eos_token_id]))
packed_labels.append(torch.tensor([-100]))
offset += 1
# 移除最后一个多余的 EOS(如果末尾不需要)
if packed_input_ids and len(packed_input_ids[-1]) == 1:
packed_input_ids.pop()
packed_labels.pop()
offset -= 1
input_ids = torch.cat(packed_input_ids, dim=0)
labels = torch.cat(packed_labels, dim=0)
total_len = input_ids.size(0)
# 构建注意力掩码:上三角因果 + 分块对角(不同样本之间屏蔽)
causal_mask = torch.tril(torch.ones(total_len, total_len)).bool()
block_mask = torch.zeros(total_len, total_len).bool()
for start, end in segment_boundaries:
block_mask[start:end+1, start:end+1] = True
attention_mask = causal_mask & block_mask # 两者同时满足
attention_mask = attention_mask.unsqueeze(0).unsqueeze(0) # (1, 1, total_len, total_len)
return {
'input_ids': input_ids.unsqueeze(0), # (1, total_len)
'labels': labels.unsqueeze(0),
'attention_mask': attention_mask
}
说明:
-
函数在每个样本后插入一个 EOS token,其 label 设为 -100,不参与损失计算。
-
注意力掩码结合了因果掩码和分块对角掩码,确保不同样本间互不干扰。
-
这种 packing 方式可提高训练效率,减少 padding 浪费。
编写脚本,从大规模jsonl文件中随机采样平衡各任务的数据¶
假设 jsonl 文件中每行是一个 JSON 对象,包含 task_type 字段。我们需要按任务类型分层采样,每类采样固定数量或比例,输出平衡后的数据。
import json
import random
from collections import defaultdict
def balanced_sampling(input_jsonl: str, output_jsonl: str, target_per_task: int = 1000, seed=42):
"""
从大 jsonl 文件中按 task_type 分层随机采样,每类采样 target_per_task 条,
不足则全取。输出到新 jsonl 文件。
"""
random.seed(seed)
task_samples = defaultdict(list)
# 第一遍:按任务分类存储(只存行号或直接存 json 字符串以省内存?这里假设文件可装入内存)
with open(input_jsonl, 'r', encoding='utf-8') as f:
for line in f:
try:
obj = json.loads(line)
task_type = obj.get('task_type', 'unknown')
task_samples[task_type].append(line)
except json.JSONDecodeError:
continue
with open(output_jsonl, 'w', encoding='utf-8') as fout:
for task, lines in task_samples.items():
if len(lines) <= target_per_task:
sampled = lines
else:
sampled = random.sample(lines, target_per_task)
for line in sampled:
fout.write(line)
print(f"采样完成,输出到 {output_jsonl}")
说明:
-
为避免内存爆炸,如果文件非常大,可以采用两遍扫描:第一遍只统计每个 task 的行号,第二遍随机读取指定行。
-
上述简单版本适用于文件可全部加载到内存的场景。
使用HuggingFace datasets库加载并预处理SFT数据,包含tokenization¶
此示例展示如何加载 jsonl 文件,将 Alpaca 格式转换为 ChatML,并进行 tokenization,最后生成input_ids 和 labels。
from datasets import Dataset
from transformers import AutoTokenizer
def preprocess_sft_data(file_path: str, tokenizer: AutoTokenizer, max_length: int = 2048):
# 定义 ChatML 格式化函数
def format_chatml(example):
instruction = example['instruction']
input_text = example.get('input', '')
output = example['output']
if input_text:
user_content = f"{instruction}\n{input_text}"
else:
user_content = instruction
# 构建 ChatML 文本
text = f"<|im_start|>user\n{user_content}<|im_end|>\n<|im_start|>assistant\n{output}<|im_end|>"
return {'text': text}
# 加载数据集
dataset = Dataset.from_json(file_path)
dataset = dataset.map(format_chatml)
# tokenize 函数
def tokenize_function(examples):
tokenized = tokenizer(
examples['text'],
truncation=True,
max_length=max_length,
padding=False,
return_tensors=None
)
# 构造 labels:复制 input_ids,然后将 prompt 部分 mask 掉
labels = []
for i, ids in enumerate(tokenized['input_ids']):
# 找到 assistant 部分开始的位置(例如 <|im_start|>assistant 后的第一个 token)
# 简单做法:利用 tokenizer 找到 'assistant' 的位置,但这里提供一种通用方法
# 我们可以先 tokenize "assistant" 得到它的 token ids,但为简化,这里假设
# 整个序列中 assistant 部分在最后一个 <|im_start|>assistant 之后。
# 更精确的实现需要解析文本,此处略。
# 下面的示例仅将所有非 -100 的设为真实 token,但在实际中应精确 mask。
lab = ids.copy()
# 暂时将全部设为 -100,然后根据规则取消 assistant 部分的 mask
# 这只是一个示意,正确实现应解析 text。
labels.append(lab)
tokenized['labels'] = labels
return tokenized
tokenized_dataset = dataset.map(tokenize_function, batched=True, remove_columns=dataset.column_names)
return tokenized_dataset
说明:
-
精确的 loss masking 需要识别
<|im_start|>assistant的位置,可使用tokenizer的encode找到分隔 token 的索引,然后对索引后的部分保留 label,前面置为 -100。 -
上述代码仅提供了框架,实际使用时应根据模型和模板进行完善。
实现计算SFT模型输出困惑度(perplexity)的代码¶
给定一个预训练模型和 tokenized 的测试集,计算平均困惑度。
import torch
import math
from torch.utils.data import DataLoader
def compute_perplexity(model, dataloader: DataLoader, device: str = 'cuda') -> float:
model.eval()
total_loss = 0.0
total_tokens = 0
with torch.no_grad():
for batch in dataloader:
input_ids = batch['input_ids'].to(device)
labels = batch['labels'].to(device)
attention_mask = batch.get('attention_mask', None)
if attention_mask is not None:
attention_mask = attention_mask.to(device)
outputs = model(
input_ids=input_ids,
attention_mask=attention_mask,
labels=labels
)
loss = outputs.loss # 模型内部已对 labels 为 -100 的部分忽略
# 计算有效 token 数
valid_tokens = (labels != -100).sum().item()
total_loss += loss.item() * valid_tokens
total_tokens += valid_tokens
avg_loss = total_loss / total_tokens
ppl = math.exp(avg_loss)
return ppl
说明:
-
依赖模型的
loss输出,它自动根据labels计算交叉熵并忽略-100。 -
困惑度是平均损失的指数。
编写函数检测两个字符串序列的最长公共子串,用于数据去污染¶
该函数用于比较训练样本与评测样本,找出最长的公共子串,当超过一定长度时可判定为污染。
def longest_common_substring(s1: str, s2: str) -> str:
"""返回两个字符串的最长公共子串(区分大小写)"""
if not s1 or not s2:
return ""
m, n = len(s1), len(s2)
dp = [[0] * (n + 1) for _ in range(m + 1)]
max_len, end_pos = 0, 0
for i in range(1, m + 1):
for j in range(1, n + 1):
if s1[i - 1] == s2[j - 1]:
dp[i][j] = dp[i - 1][j - 1] + 1
if dp[i][j] > max_len:
max_len = dp[i][j]
end_pos = i
return s1[end_pos - max_len: end_pos]
def detect_contamination(train_sample: str, eval_sample: str, threshold: int = 50) -> bool:
"""
如果最长公共子串长度 >= threshold,则认为存在污染。
"""
lcs = longest_common_substring(train_sample, eval_sample)
return len(lcs) >= threshold
说明:
- 对于长文本,该 O(n*m) 算法可能较慢,可采用后缀自动机或二分+哈希优化,但基本方法足以说明。
实现基于MinHash的文本去重代码¶
MinHash 可用于快速估计两个集合的 Jaccard 相似度,常用于海量文本的近似去重。
import re
import hashlib
from collections import defaultdict
class MinHash:
def __init__(self, num_perm: int = 128, seed: int = 42):
self.num_perm = num_perm
# 生成随机排列参数 (a, b) 用于模拟哈希函数 h(x) = (a*x + b) mod prime
self.seed = seed
self.prime = 2**61 - 1
self.a = self._random_coeffs()
self.b = self._random_coeffs()
def _random_coeffs(self):
import random
random.seed(self.seed)
return [random.randint(1, self.prime - 1) for _ in range(self.num_perm)]
def _hash_func(self, x, i):
return (self.a[i] * x + self.b[i]) % self.prime
def compute_signature(self, tokens: set) -> list:
# tokens 为字符串集合,将其哈希为整数
token_hashes = [hash(t) & 0xffffffff for t in tokens] # 转为 32位整数
if not token_hashes:
return [float('inf')] * self.num_perm
signature = []
for i in range(self.num_perm):
min_val = min(self._hash_func(h, i) for h in token_hashes)
signature.append(min_val)
return signature
@staticmethod
def similarity(sig1: list, sig2: list) -> float:
"""估计 Jaccard 相似度"""
return sum(x == y for x, y in zip(sig1, sig2)) / len(sig1)
def deduplicate_texts(texts: list[str], threshold: float = 0.8, ngram_size: int = 3) -> list[str]:
"""基于 MinHash 去重,返回去重后的文本列表"""
mh = MinHash(num_perm=128)
signatures = []
for text in texts:
# 提取 n-gram 集合
words = text.split()
ngrams = set()
for i in range(len(words) - ngram_size + 1):
ngrams.add(' '.join(words[i:i+ngram_size]))
signatures.append(mh.compute_signature(ngrams))
keep = []
for i, text in enumerate(texts):
is_dup = False
for j in keep:
if mh.similarity(signatures[i], signatures[j]) >= threshold:
is_dup = True
break
if not is_dup:
keep.append(i)
return [texts[i] for i in keep]
说明:
-
使用 n-gram 将文本转换为 token 集合,再计算 MinHash 签名。
-
通过比较签名相似度(估计 Jaccard)来判断重复,避免了逐对文本的 O(n^2) 比较。
编写一个简单的Self-Instruct指令生成循环,调用OpenAI API¶
该循环用种子指令和 few-shot 示例,调用 GPT API 生成新的指令及其回答,并进行简单的多样性过滤。
import openai
import random
import json
openai.api_key = "your-api-key"
def generate_new_instructions(seed_tasks: list, num_generate: int = 100, model="gpt-4"):
generated = []
for _ in range(num_generate):
# 随机选取 6 个种子任务作为 few-shot 示例
few_shot = random.sample(seed_tasks, min(6, len(seed_tasks)))
# 构造 prompt
prompt = "你是一个智能助手,请生成一个全新的、不同于以下示例的指令,并为该指令提供一个高质量的回答。\n\n"
for task in few_shot:
prompt += f"指令: {task['instruction']}\n回答: {task['output']}\n\n"
prompt += "现在,请生成一个全新的指令和回答。\n指令:"
response = openai.ChatCompletion.create(
model=model,
messages=[{"role": "user", "content": prompt}],
temperature=0.8,
max_tokens=256
)
content = response.choices[0].message.content.strip()
# 简单解析:假设格式为 "指令: ...\n回答: ..."
parts = content.split('\n回答:', 1) # 可能包含多种格式,简化
if len(parts) == 2:
instruction = parts[0].replace('指令:', '').strip()
output = parts[1].strip()
# 多样性过滤:与已有指令计算相似度(略),这里仅去重
if not any(instruction == t['instruction'] for t in generated):
generated.append({'instruction': instruction, 'output': output})
return generated
说明:
-
实际 Self-Instruct 流程会包含更复杂的解析和过滤(如 ROUGE-L 相似度检查),此处仅给出基本框架。
-
需要安装 openai 库并配置 API 密钥。
如何实现一个流式数据处理管道,用于海量SFT数据清洗?¶
对于无法完全装入内存的海量数据,我们使用生成器(generator)和管道模式,逐步读取、处理、过滤并写入结果。
import json
import gzip
from typing import Iterator, Callable
def stream_jsonl(input_files: list[str]) -> Iterator[dict]:
"""逐个读取 jsonl(或 .jsonl.gz)文件,yield 每条记录"""
for file_path in input_files:
opener = gzip.open if file_path.endswith('.gz') else open
with opener(file_path, 'rt', encoding='utf-8') as f:
for line in f:
line = line.strip()
if not line:
continue
try:
obj = json.loads(line)
yield obj
except json.JSONDecodeError:
continue
def clean_pipeline(input_files: list[str], output_file: str, filter_func: Callable, transform_func=None):
"""
流式处理:读取 -> 过滤 -> 转换 -> 写入。
filter_func(obj) 返回 True 保留;transform_func(obj) 可修改对象。
"""
with open(output_file, 'w', encoding='utf-8') as fout:
for obj in stream_jsonl(input_files):
if filter_func(obj):
if transform_func:
obj = transform_func(obj)
fout.write(json.dumps(obj, ensure_ascii=False) + '\n')
print(f"流式处理完成,输出到 {output_file}")
# 示例过滤函数:保留指令长度>10且回答长度>5的样本
def length_filter(obj):
return len(obj.get('instruction', '')) > 10 and len(obj.get('output', '')) > 5
说明:
-
使用
Iterator模式逐行读取,内存占用极小。 -
可根据需求组合多个过滤和转换函数,实现复杂的清洗流水线。
写一个评估函数,计算模型在IFEval指令遵循准确率¶
IFEval 评估的是模型是否严格遵循指令中的可验证约束,例如字数限制、包含特定关键词等。这里提供一个简化版框架,展示如何按约束检查并计算准确率。
import re
from typing import List, Dict, Callable
def ifeval_evaluate(model_responses: List[str], constraints: List[Dict]) -> float:
"""
model_responses: 模型生成的回答列表
constraints: 每条指令对应的约束列表,每个约束是一个字典,包含 'type' 和 'params'。
返回严格指令遵循准确率(所有约束都满足的样本比例)。
"""
assert len(model_responses) == len(constraints), "样本数不匹配"
total = len(model_responses)
strict_correct = 0
for response, constraint_list in zip(model_responses, constraints):
all_pass = True
for con in constraint_list:
check_func = CONSTRAINT_CHECKERS.get(con['type'])
if not check_func:
continue
if not check_func(response, con.get('params', {})):
all_pass = False
break
if all_pass:
strict_correct += 1
return strict_correct / total
# 预定义的约束检查函数
CONSTRAINT_CHECKERS = {}
def register_checker(type_name: str):
def decorator(func):
CONSTRAINT_CHECKERS[type_name] = func
return func
return decorator
@register_checker('word_count_between')
def check_word_count(text: str, params: dict) -> bool:
"""字数在 [min, max] 之间(中英文通用简单分词)"""
# 简化:按空白字符分词
words = text.split()
min_w = params.get('min', 0)
max_w = params.get('max', float('inf'))
return min_w <= len(words) <= max_w
@register_checker('must_contain')
def check_must_contain(text: str, params: dict) -> bool:
"""文本中必须包含指定子串"""
substring = params['substring']
return substring in text
@register_checker('forbidden_words')
def check_forbidden_words(text: str, params: dict) -> bool:
"""不能包含指定词语列表中的任何一个"""
forbidden = params.get('words', [])
for word in forbidden:
if word in text:
return False
return True
@register_checker('json_format')
def check_json_format(text: str, params: dict = None) -> bool:
"""检查文本是否是有效的 JSON(去除前后空白后)"""
import json
try:
json.loads(text.strip())
return True
except:
return False
说明:
-
用户可根据 IFEval 的约束类型扩展
checker函数。 -
该函数直接返回所有约束都通过的“严格准确率”。
-
调用时,需要事先为每条测试指令定义好约束列表。
实现GPT-as-judge的评估脚本,包含比较两个模型回复的逻辑¶
import openai
import json
import random
from typing import List, Dict
openai.api_key = "your-api-key"
def gpt_judge_compare(prompt: str, response_a: str, response_b: str,
model="gpt-4", temperature=0.0) -> Dict:
"""
使用GPT模型比较两个回复,返回评判结果。
输出格式: {'winner': 'A'/'B'/'tie', 'reason': '...'}
"""
# 随机交换顺序以避免位置偏差
if random.random() < 0.5:
first, second = "A", "B"
resp1, resp2 = response_a, response_b
else:
first, second = "B", "A"
resp1, resp2 = response_b, response_a
judge_prompt = f"""请作为一个公正的评判者,比较以下两个AI助手对用户问题的回复。
用户问题: {prompt}
助手{first}的回复: {resp1}
助手{second}的回复: {resp2}
请从准确性、有用性、流畅度和完整性等方面综合评判。
只输出JSON格式: {{"winner": "{first}"或"{second}"或"tie", "reason": "简要理由"}}
"""
response = openai.ChatCompletion.create(
model=model,
messages=[{"role": "user", "content": judge_prompt}],
temperature=temperature,
max_tokens=256
)
content = response.choices[0].message.content.strip()
try:
result = json.loads(content)
# 如果因为顺序交换,需要将winner映射回原始的A/B
if result['winner'] == first:
result['winner'] = 'A'
elif result['winner'] == second:
result['winner'] = 'B'
except json.JSONDecodeError:
result = {'winner': 'error', 'reason': content}
return result
def batch_compare(prompts: List[str], responses_a: List[str], responses_b: List[str]) -> Dict:
"""批量比较,统计胜率"""
stats = {'A_wins': 0, 'B_wins': 0, 'ties': 0}
for p, ra, rb in zip(prompts, responses_a, responses_b):
judge = gpt_judge_compare(p, ra, rb)
if judge['winner'] == 'A':
stats['A_wins'] += 1
elif judge['winner'] == 'B':
stats['B_wins'] += 1
else:
stats['ties'] += 1
return stats
编写代码可视化SFT训练过程中的attention分布¶
可以使用HuggingFace的output_attentions=True获取注意力矩阵,并使用matplotlib绘制热力图。
import torch
import matplotlib.pyplot as plt
import seaborn as sns
from transformers import AutoModelForCausalLM, AutoTokenizer
def visualize_attention(model_name: str, text: str, layer_idx: int = 0, head_idx: int = 0):
tokenizer = AutoTokenizer.from_pretrained(model_name)
model = AutoModelForCausalLM.from_pretrained(model_name, output_attentions=True)
inputs = tokenizer(text, return_tensors='pt')
with torch.no_grad():
outputs = model(**inputs)
# attentions 是 tuple,每层形状 (batch, num_heads, seq_len, seq_len)
attn_weights = outputs.attentions[layer_idx][0, head_idx].numpy()
tokens = tokenizer.convert_ids_to_tokens(inputs['input_ids'][0])
plt.figure(figsize=(10, 8))
sns.heatmap(attn_weights, xticklabels=tokens, yticklabels=tokens, cmap='Blues')
plt.title(f'Attention Weights (Layer {layer_idx}, Head {head_idx})')
plt.tight_layout()
plt.show()
# 使用示例
# visualize_attention('gpt2', '今天天气很好', layer_idx=0, head_idx=0)
实现一个简单的梯度累积训练循环(不使用Trainer)¶
手动实现包含梯度累积的PyTorch训练循环。
import torch
from torch.optim import AdamW
from tqdm import tqdm
def train_with_gradient_accumulation(model, train_dataloader, num_epochs=3,
accumulation_steps=4, lr=5e-5):
optimizer = AdamW(model.parameters(), lr=lr)
model.train()
global_step = 0
for epoch in range(num_epochs):
total_loss = 0
optimizer.zero_grad()
progress_bar = tqdm(train_dataloader, desc=f'Epoch {epoch+1}')
for step, batch in enumerate(progress_bar):
outputs = model(**batch)
loss = outputs.loss / accumulation_steps # 归一化损失
loss.backward()
if (step + 1) % accumulation_steps == 0:
optimizer.step()
optimizer.zero_grad()
global_step += 1
total_loss += loss.item() * accumulation_steps
progress_bar.set_postfix({'loss': total_loss / (step + 1)})
使用DeepSpeed配置文件启动SFT训练,写出关键参数¶
DeepSpeed配置文件 ds_config.json 示例(ZeRO-2 + CPU offload):
{
"train_batch_size": 16,
"gradient_accumulation_steps": 4,
"optimizer": {
"type": "AdamW",
"params": {
"lr": 2e-5,
"weight_decay": 0.1
}
},
"fp16": {
"enabled": false
},
"bf16": {
"enabled": true
},
"zero_optimization": {
"stage": 2,
"offload_optimizer": {
"device": "cpu",
"pin_memory": true
},
"allgather_partitions": true,
"allgather_bucket_size": 5e8,
"overlap_comm": true,
"reduce_scatter": true,
"reduce_bucket_size": 5e8,
"contiguous_gradients": true
},
"gradient_clipping": 1.0,
"steps_per_print": 10
}
启动命令:
如何用HuggingFace PEFT库加载一个QLoRA模型并进行推理?¶
from transformers import AutoModelForCausalLM, AutoTokenizer, BitsAndBytesConfig
from peft import PeftModel, PeftConfig
# 1. 加载基座模型(量化)
bnb_config = BitsAndBytesConfig(
load_in_4bit=True,
bnb_4bit_quant_type="nf4",
bnb_4bit_compute_dtype=torch.bfloat16,
bnb_4bit_use_double_quant=True
)
model = AutoModelForCausalLM.from_pretrained(
"meta-llama/Llama-2-7b-hf",
quantization_config=bnb_config,
device_map="auto"
)
tokenizer = AutoTokenizer.from_pretrained("meta-llama/Llama-2-7b-hf")
# 2. 加载LoRA适配器
peft_model = PeftModel.from_pretrained(model, "path/to/lora-adapter")
# 3. 推理
inputs = tokenizer("你是谁?", return_tensors="pt").to("cuda")
outputs = peft_model.generate(**inputs, max_new_tokens=100)
print(tokenizer.decode(outputs[0], skip_special_tokens=True))
编写代码统计SFT数据集中指令的长度分布和回复的长度分布¶
import json
import matplotlib.pyplot as plt
from transformers import AutoTokenizer
def analyze_length_distribution(jsonl_file: str, tokenizer_name: str = 'gpt2'):
tokenizer = AutoTokenizer.from_pretrained(tokenizer_name)
instr_lens = []
resp_lens = []
with open(jsonl_file, 'r', encoding='utf-8') as f:
for line in f:
obj = json.loads(line)
instr = obj.get('instruction', '')
resp = obj.get('output', '')
instr_lens.append(len(tokenizer.encode(instr)))
resp_lens.append(len(tokenizer.encode(resp)))
fig, (ax1, ax2) = plt.subplots(1, 2, figsize=(12, 5))
ax1.hist(instr_lens, bins=50, alpha=0.7)
ax1.set_title('Instruction Length Distribution')
ax2.hist(resp_lens, bins=50, alpha=0.7, color='orange')
ax2.set_title('Response Length Distribution')
plt.show()
return instr_lens, resp_lens
实现一个简单的早期停止(Early Stopping)类¶
class EarlyStopping:
def __init__(self, patience=3, min_delta=0.0, mode='min'):
self.patience = patience
self.min_delta = min_delta
self.mode = mode
self.counter = 0
self.best_score = None
self.early_stop = False
def __call__(self, current_score):
if self.best_score is None:
self.best_score = current_score
elif (self.mode == 'min' and current_score > self.best_score - self.min_delta) or \
(self.mode == 'max' and current_score < self.best_score + self.min_delta):
self.counter += 1
if self.counter >= self.patience:
self.early_stop = True
else:
self.best_score = current_score
self.counter = 0
return self.early_stop
写代码计算两个SFT模型在同一提示下的输出KL散度¶
import torch
import torch.nn.functional as F
def compute_kl_divergence(model1, model2, tokenizer, prompt, max_len=50):
"""计算两个模型对同一prompt生成序列时的KL散度"""
inputs = tokenizer(prompt, return_tensors='pt')
# 使用模型1生成序列
gen_out = model1.generate(**inputs, max_new_tokens=max_len, do_sample=False,
output_scores=True, return_dict_in_generate=True)
generated_ids = gen_out.sequences[0]
# 获取两个模型在每个生成token上的logits
logits1 = []
logits2 = []
with torch.no_grad():
past_kv1 = None
past_kv2 = None
input_ids = inputs.input_ids
for _ in range(generated_ids.size(-1) - input_ids.size(-1)):
out1 = model1(input_ids, past_key_values=past_kv1, use_cache=True)
out2 = model2(input_ids, past_key_values=past_kv2, use_cache=True)
logits1.append(out1.logits[:, -1, :])
logits2.append(out2.logits[:, -1, :])
past_kv1 = out1.past_key_values
past_kv2 = out2.past_key_values
# 下一步输入是当前生成的token
input_ids = generated_ids[:, len(inputs.input_ids[0]):len(inputs.input_ids[0])+1+_]
logits1 = torch.cat(logits1, dim=0)
logits2 = torch.cat(logits2, dim=0)
p = F.softmax(logits1, dim=-1)
q = F.softmax(logits2, dim=-1)
kl = F.kl_div(F.log_softmax(logits1, dim=-1), q, reduction='batchmean')
return kl.item()
实现基于困惑度的数据过滤脚本,过滤高PPL的样本¶
import torch
from transformers import AutoModelForCausalLM, AutoTokenizer
from tqdm import tqdm
def filter_by_ppl(input_file, output_file, model_name='gpt2', threshold=1000.0):
tokenizer = AutoTokenizer.from_pretrained(model_name)
model = AutoModelForCausalLM.from_pretrained(model_name).cuda()
model.eval()
kept = 0
total = 0
with open(input_file, 'r') as fin, open(output_file, 'w') as fout:
for line in tqdm(fin):
total += 1
data = json.loads(line)
text = data['output'] # 仅评估回复部分的PPL
inputs = tokenizer(text, return_tensors='pt').to('cuda')
with torch.no_grad():
outputs = model(**inputs, labels=inputs['input_ids'])
loss = outputs.loss
ppl = torch.exp(loss).item()
if ppl < threshold:
fout.write(line)
kept += 1
print(f'保留 {kept}/{total} 条, 过滤率 {1-kept/total:.2%}')
编写函数,对SFT数据中的指令进行“进化”操作(简单实现Evol-Instruct)¶
import openai
import json
def evolve_instruction(seed_instruction: str, evolution_type: str = "depth") -> str:
"""
evolution_type: 'depth' 增加复杂度/约束, 'breadth' 主题迁移/改写
"""
if evolution_type == "depth":
prompt = f"""请将以下指令改写得更加复杂和具有挑战性,增加更多限制条件或需要多步推理。
原指令: {seed_instruction}
输出只包含改写后的指令,不需要任何额外文本。"""
else:
prompt = f"""请将以下指令的主题或领域进行迁移,保持任务结构不变,但用完全不同的内容替换。
原指令: {seed_instruction}
输出只包含新指令。"""
response = openai.ChatCompletion.create(
model="gpt-4",
messages=[{"role": "user", "content": prompt}],
temperature=0.7,
max_tokens=200
)
return response.choices[0].message.content.strip()
如何用vLLM部署一个SFT后的模型并实现一个简单的对话客户端?¶
部署命令(服务端):
python -m vllm.entrypoints.openai.api_server \
--model /path/to/sft-model \
--gpu-memory-utilization 0.9 \
--max-model-len 4096
客户端示例(Python):
import requests
def chat(prompt, history=None):
url = "http://localhost:8000/v1/chat/completions"
messages = [{"role": "system", "content": "你是一个有帮助的助手。"}]
if history:
messages.extend(history)
messages.append({"role": "user", "content": prompt})
resp = requests.post(url, json={
"model": "sft-model",
"messages": messages,
"max_tokens": 256,
"temperature": 0.7
})
return resp.json()['choices'][0]['message']['content']
编写脚本自动化进行SFT模型的A/B测试,返回胜率统计¶
假设已有API endpoints。
import requests
import random
def ab_test(prompts, model_a_url, model_b_url, judge_model="gpt-4"):
wins = {'A': 0, 'B': 0, 'tie': 0}
for prompt in prompts:
resp_a = requests.post(model_a_url, json={"prompt": prompt}).json()['response']
resp_b = requests.post(model_b_url, json={"prompt": prompt}).json()['response']
# 随机交换位置以调用GPT评判
if random.random() < 0.5:
winner = gpt_judge_compare(prompt, resp_a, resp_b)['winner']
else:
winner = gpt_judge_compare(prompt, resp_b, resp_a)['winner']
winner = 'A' if winner == 'B' else 'B' if winner == 'A' else 'tie'
wins[winner] += 1
total = len(prompts)
print(f"A胜率: {wins['A']/total:.2%}, B胜率: {wins['B']/total:.2%}, 平局: {wins['tie']/total:.2%}")
实现一个对话历史管理器,用于多轮SFT数据构建¶
class DialogueManager:
def __init__(self, system_prompt=None):
self.system_prompt = system_prompt
self.turns = []
def add_user(self, content):
self.turns.append(('user', content))
def add_assistant(self, content):
self.turns.append(('assistant', content))
def format_chatml(self) -> str:
result = ""
if self.system_prompt:
result += f"<|im_start|>system\n{self.system_prompt}<|im_end|>\n"
for role, content in self.turns:
result += f"<|im_start|>{role}\n{content}<|im_end|>\n"
return result
def extract_training_sample(self) -> dict:
# 返回最后一个 assistant 回复作为目标,之前的内容作为上下文
# 实际中更复杂,这里简化
return self.format_chatml()
使用transformers的Trainer自定义loss函数,实现加权的SFT损失¶
可以通过重写Trainer的compute_loss方法来实现。
from transformers import Trainer
import torch
class WeightedTrainer(Trainer):
def __init__(self, task_weights=None, *args, **kwargs):
super().__init__(*args, **kwargs)
self.task_weights = task_weights # dict: task_type -> weight
def compute_loss(self, model, inputs, return_outputs=False):
labels = inputs.get("labels")
outputs = model(**inputs)
logits = outputs.logits
# 计算每个token的交叉熵
loss_fct = torch.nn.CrossEntropyLoss(reduction='none', ignore_index=-100)
loss_per_token = loss_fct(logits.view(-1, logits.size(-1)), labels.view(-1))
# 根据task_type赋予不同权重(需在数据中传入task_type并扩展到token维度)
# 简单示例:若数据中包含 task_weight 字段
weights = inputs.get('task_weight', torch.ones_like(labels)).float()
weights = weights.view(-1)
weighted_loss = (loss_per_token * weights).sum() / (weights.sum() + 1e-8)
return (weighted_loss, outputs) if return_outputs else weighted_loss
使用时在数据中传入task_weight字段,例如通过DataCollator将每类的权重扩展为与labels相同形状的张量。
以上实现涵盖了要求的21至34题的核心代码,可在实际项目中直接使用或根据需求微调。