跳转至

八:手撕代码与项目深挖

实现一个将对话数据转换为训练格式的Data Collator,包含loss masking

这个 Data Collator 接收一个 batch 的对话样本,每个样本包含 input_idslabels(已经对 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 格式通常包含 instructioninput(可选)和 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_idslabels

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 的位置,可使用 tokenizerencode 找到分隔 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
}

启动命令:

deepspeed --num_gpus=4 train.py --deepspeed ds_config.json

如何用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题的核心代码,可在实际项目中直接使用或根据需求微调。