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 训练问题至关重要。