跳转至

util.py — 序列打包/解包工具

文件路径

verl/models/mcore/util.py

文件概述

提供序列打包 (packing) 和解包 (unpacking) 的工具函数。Megatron-Core 支持两种数据格式:THD(Token-Head-Dimension,无 padding)和 BSHD(Batch-Sequence-Head-Dimension,有 padding)。这个文件处理两种格式的预处理和后处理,包括 Context Parallel (CP) 支持。

关键代码讲解

THD 格式预处理 -- preprocess_packed_seqs()

def preprocess_packed_seqs(input_ids, attention_mask, pre_process=True):
    """将 batch 中的多个序列打包为连续的 token 序列(THD 格式)"""
    seqlens_in_batch = attention_mask.sum(dim=-1, dtype=torch.int32)
    # 对齐到 TP * CP * 2
    align_size = tp_size * cp_size * 2 if cp_size > 1 else tp_size
    pad_size = (align_size - seqlens_in_batch % align_size) % align_size

    # 构建 cu_seqlens(累积序列长度)
    cu_seqlens_padded[1:] = torch.cumsum(seqlens_in_batch_padded, dim=0)

    # CP 切分:每个 GPU 获取序列的前半和后半各一块
    if cp_size > 1:
        input_ids_rmpad[start_idx : start_idx + half_seqlen] = d[half_seqlen * cp_rank : ...]
        # 第二块取序列末尾
        input_ids_rmpad[start_idx + half_seqlen : ...] = d[remain_start:remain_end]

    packed_seq_params = PackedSeqParams(
        qkv_format="thd",
        cu_seqlens_q=cu_seqlens_padded,
        max_seqlen_q=max_seqlen_in_batch,
    )
    return input_ids_rmpad.unsqueeze(0), packed_seq_params

CP 切分策略:Context Parallel 将序列切成 CP*2 块,GPU0 拿第一块和最后一块,GPU1 拿第二块和倒数第二块...这样做是为了在因果注意力中实现负载均衡。

THD 格式后处理 -- postprocess_packed_seqs()

def postprocess_packed_seqs(output, packed_seq_params, attention_mask, batch_size, seq_len):
    """将打包的输出解包回原始 batch 格式"""
    if cp_size > 1:
        # all_gather 跨 CP 组收集输出
        output_list = [torch.empty_like(output) for _ in range(cp_size)]
        torch.distributed.all_gather(output_list, output.detach(), group=mpu.get_context_parallel_group())
        # 重新排列回原始序列顺序

BSHD 格式预处理 -- preprocess_bshd()

def preprocess_bshd(input_ids, attention_mask, position_ids, sequence_parallel=False):
    """去除左 padding,右对齐序列"""
    seq_len = seq_lens.max().item()
    if sequence_parallel:
        # SP 需要序列长度对齐到 TP 大小
        pad_size = (sp_world_size - seq_len % sp_world_size) % sp_world_size
        seq_len = seq_len + pad_size
    # 将每个序列的有效 token 左对齐
    new_input_ids[i, : seq_lens[i]] = input_ids[i, attention_mask[i]]

无 padding 版本(Nested Tensor 输入)

def preprocess_thd_no_padding(input_ids, pre_process=True, need_roll=False):
    """输入为 nested tensor 的 THD 预处理"""
    seqlens_in_batch = input_ids.offsets().diff()  # 从 nested tensor 获取序列长度
    # ... 与 preprocess_packed_seqs 类似

def postprocess_thd_no_padding(output, packed_seq_params, input_ids, batch_size):
    """输出转回 nested tensor"""
    output_new = []
    for i in range(batch_size):
        output_new.append(output[0][start_idx : start_idx + s])
    return torch.nested.as_nested_tensor(output_new, layout=torch.jagged)

融合 kernel 输出后处理

def postprocess_packed_seqs_for_dict_output(labels_mask, output, packed_seq_params, ...):
    """对融合 kernel 输出的 log_probs 和 entropy 分别解包"""
    output.log_probs = output.log_probs.masked_fill(~labels_mask, 0.0)
    ret["entropy"] = postprocess_packed_seqs(output.entropy, ...)
    ret["log_probs"] = postprocess_packed_seqs(output.log_probs, ...)
    return ret

核心函数列表

函数名 作用
preprocess_packed_seqs() THD 打包预处理(含 CP)
postprocess_packed_seqs() THD 解包后处理(含 CP)
preprocess_bshd() BSHD 去 padding 预处理
postprocess_bshd() BSHD 恢复 padding 后处理
preprocess_thd_no_padding() Nested tensor THD 预处理
postprocess_thd_no_padding() Nested tensor THD 后处理
preprocess_bshd_no_padding() Nested tensor BSHD 预处理
postprocess_bshd_no_padding() Nested tensor BSHD 后处理
postprocess_packed_seqs_for_dict_output() 融合 kernel 字典输出后处理

与其他模块的关系

  • 被 model_forward.py 和 model_forward_fused.py 调用
  • 与 Megatron-Core 的 PackedSeqParams 配合使用

小结

这个文件处理了 mcore 训练中最繁琐的数据格式转换。THD 格式通过序列打包避免了 padding 浪费,但需要复杂的预处理(对齐、CP 切分)和后处理(解包、all_gather)。理解这些转换对于调试 mcore 训练问题至关重要。